diff --git a/.fern/metadata.json b/.fern/metadata.json index f5b9f8fd..10d57bdd 100644 --- a/.fern/metadata.json +++ b/.fern/metadata.json @@ -12,7 +12,7 @@ "enable-wire-tests": true, "runtime-version": true }, - "originGitCommit": "e252995bcdb25f3e12d46ae342a2b93d0c1085f9", + "originGitCommit": "9dcbd3d5e8d36420319e1b33988613ad2cb2be10", "originGitCommitIsDirty": true, "invokedBy": "manual", "sdkVersion": "0.11.0" diff --git a/.fernignore b/.fernignore index ef6a8355..550252f0 100644 --- a/.fernignore +++ b/.fernignore @@ -21,6 +21,11 @@ src/main/java/com/deepgram/core/ClientOptions.java # Transport abstraction (pluggable transport for SageMaker, etc.) src/main/java/com/deepgram/core/transport/ +# Hand-written exception that carries a Listen V1 server Error to the generic onError handler +# (thrown by the frozen listen v1 V1WebSocketClient below when no onErrorMessage handler is set). +# No Fern equivalent; the generator would delete it. +src/main/java/com/deepgram/resources/listen/v1/websocket/ListenV1ErrorException.java + # 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. diff --git a/AGENTS.md b/AGENTS.md index bcef6fd2..7f6717e6 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -26,6 +26,7 @@ Current permanently frozen files: - `src/main/java/com/deepgram/DeepgramClient.java`, `src/main/java/com/deepgram/AsyncDeepgramClient.java`, `src/main/java/com/deepgram/DeepgramClientBuilder.java`, `src/main/java/com/deepgram/AsyncDeepgramClientBuilder.java` - custom wrapper entrypoints that add Bearer auth, session ID support, and custom transport behavior on top of Fern's generated API client - `src/main/java/com/deepgram/core/transport/` - hand-written transport abstraction +- `src/main/java/com/deepgram/resources/listen/v1/websocket/ListenV1ErrorException.java` - hand-written exception carrying a Listen V1 server `Error` (`ListenV1Error`) to the generic `onError` handler when no `onErrorMessage` handler is registered; thrown by the frozen listen v1 `V1WebSocketClient`. No Fern equivalent - `build.gradle`, `settings.gradle`, `gradle/`, `gradlew`, `gradlew.bat`, `pom.xml`, `Makefile` - build and project configuration - `README.md`, `CHANGELOG.md`, `CONTRIBUTING.md`, `LICENSE`, `docs/` - docs - `src/test/` - manually maintained tests @@ -47,10 +48,10 @@ 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. Fern 4.22.1 tracks and closes child WebSockets before shutting down SDK-owned OkHttp resources; the permanent custom client wrappers inherit it, so do not add wrapper-local lifecycle code and retain Fern's lifecycle implementation during reconciliation. +- `src/main/java/com/deepgram/core/ClientOptions.java` - preserves release-please version markers, correct SDK header constants, and the accurate `close()` lifecycle Javadoc 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. Fern 4.22.1 tracks and closes child WebSockets before shutting down SDK-owned OkHttp resources; the permanent custom client wrappers inherit it, so do not add wrapper-local lifecycle code and retain Fern's lifecycle implementation during reconciliation. - `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/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). +- `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. Listen V1 also forwards a parsed server `Error` to `onError` as a `ListenV1ErrorException` (carrying the typed `ListenV1Error`) when no `onErrorMessage` handler is registered, preserving the prior generic error-handler behavior. 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). - Union default-variant fix on the agent listen-provider unions: `src/main/java/com/deepgram/resources/agent/v1/types/AgentV1UpdateListenListenProvider.java`, `src/main/java/com/deepgram/resources/agent/v1/types/AgentV1SettingsAgentListenProvider.java`, `src/main/java/com/deepgram/resources/agent/v1/types/AgentV1SettingsAgentContextListenProvider.java`. `version` is an optional discriminator, so a provider payload without it is valid (and is what 0.7.x emits), but Fern points `@JsonTypeInfo` `defaultImpl` at the empty-bodied `_UnknownValue`, so such a payload deserializes to an unknown variant carrying `null` — `getProvider()` returns `null` and re-serialization emits `{"provider":null}`, silently dropping the provider on the wire. Patched to `defaultImpl = V2Value` on each; guarded by `src/test/java/com/deepgram/AgentSettingsProviderDefaultTest.java`. Use the standard `.bak` swap/restore workflow during regen review; drop the patches and unfreeze once the generator stops defaulting unions to the empty `_UnknownValue` (tracked as an upstream Fern request). diff --git a/README.md b/README.md index bab5073e..3940e280 100644 --- a/README.md +++ b/README.md @@ -259,11 +259,14 @@ Stream audio for real-time speech-to-text. import com.deepgram.DeepgramClient; import com.deepgram.resources.listen.v1.types.ListenV1CloseStream; import com.deepgram.resources.listen.v1.types.ListenV1CloseStreamType; +import com.deepgram.resources.listen.v1.types.ListenV1Configure; import com.deepgram.resources.listen.v1.websocket.V1WebSocketClient; import com.deepgram.resources.listen.v1.websocket.V1ConnectOptions; import com.deepgram.types.ListenV1Model; import java.nio.file.Files; import java.nio.file.Path; +import java.util.List; +import java.util.Map; import java.util.concurrent.TimeUnit; import okio.ByteString; @@ -288,12 +291,25 @@ ws.onError(error -> { System.err.println("Error: " + error.getMessage()); }); +ws.onErrorMessage(error -> { + System.err.println("Server error: " + error.getVariant() + ": " + error.getDescription()); +}); + // Connect with options (model is required) ws.connect(V1ConnectOptions.builder() .model(ListenV1Model.NOVA3) .build()) .get(10, TimeUnit.SECONDS); +// Update Nova-3 keyterms and numerals without reconnecting. Keep the keyterm list under the +// 500-token limit: an over-limit update currently stops transcription without an Error, and +// the server closes the stream. +ws.sendConfigure(ListenV1Configure.builder() + .keyterms(List.of("Deepgram")) + .features(Map.of("numerals", true)) + .build()) + .get(5, TimeUnit.SECONDS); + ws.sendMedia(ByteString.of(audioBytes)); ws.sendCloseStream(ListenV1CloseStream.builder() .type(ListenV1CloseStreamType.CLOSE_STREAM) diff --git a/examples/listen/LiveReconfigure.java b/examples/listen/LiveReconfigure.java new file mode 100644 index 00000000..6da7fcca --- /dev/null +++ b/examples/listen/LiveReconfigure.java @@ -0,0 +1,86 @@ +import com.deepgram.DeepgramClient; +import com.deepgram.resources.listen.v1.types.ListenV1CloseStream; +import com.deepgram.resources.listen.v1.types.ListenV1CloseStreamType; +import com.deepgram.resources.listen.v1.types.ListenV1Configure; +import com.deepgram.resources.listen.v1.types.ListenV1Error; +import com.deepgram.resources.listen.v1.types.ListenV1ResultsChannelAlternativesItem; +import com.deepgram.resources.listen.v1.websocket.V1ConnectOptions; +import com.deepgram.resources.listen.v1.websocket.V1WebSocketClient; +import com.deepgram.types.ListenV1Model; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +/** + * Reconfigures an active Nova-3 Listen V1 stream without reconnecting. + * + *

Keep the keyterm list under the 500-token limit: an over-limit update currently stops transcription without an + * {@code Error}, and the server closes the stream. + * + *

Usage: {@code DEEPGRAM_API_KEY=... java LiveReconfigure} + */ +public class LiveReconfigure { + public static void main(String[] args) { + String apiKey = System.getenv("DEEPGRAM_API_KEY"); + if (apiKey == null || apiKey.isEmpty()) { + System.err.println("DEEPGRAM_API_KEY environment variable is required"); + System.exit(1); + } + + DeepgramClient client = DeepgramClient.builder().apiKey(apiKey).build(); + V1WebSocketClient wsClient = client.listen().v1().v1WebSocket(); + CountDownLatch closeLatch = new CountDownLatch(1); + + try { + wsClient.onConnected(() -> System.out.println("Connected to Deepgram")); + wsClient.onResults(result -> { + if (result.getChannel() != null + && result.getChannel().getAlternatives() != null + && !result.getChannel().getAlternatives().isEmpty()) { + ListenV1ResultsChannelAlternativesItem alternative = + result.getChannel().getAlternatives().get(0); + if (alternative.getTranscript() != null + && !alternative.getTranscript().isEmpty()) { + System.out.println(alternative.getTranscript()); + } + } + }); + wsClient.onErrorMessage(LiveReconfigure::printListenError); + wsClient.onError(error -> System.err.println("WebSocket error occurred: " + error.getMessage())); + wsClient.onDisconnected(reason -> { + System.out.println("Connection closed."); + closeLatch.countDown(); + }); + + wsClient.connect(V1ConnectOptions.builder() + .model(ListenV1Model.NOVA3) + .build()) + .get(10, TimeUnit.SECONDS); + + wsClient.sendConfigure(ListenV1Configure.builder() + .keyterms(List.of("Deepgram", "Nova-3")) + .features(Map.of("numerals", true)) + .build()) + .get(10, TimeUnit.SECONDS); + System.out.println("Sent live reconfiguration for keyterms and numerals."); + + // Stream audio after this point, for example: wsClient.sendMedia(audioChunk); + wsClient.sendCloseStream(ListenV1CloseStream.builder() + .type(ListenV1CloseStreamType.CLOSE_STREAM) + .build()) + .get(10, TimeUnit.SECONDS); + closeLatch.await(15, TimeUnit.SECONDS); + } catch (Exception e) { + System.err.println("Unable to run the live reconfiguration example: " + e.getMessage()); + } finally { + wsClient.close(); + client.close(); + } + } + + private static void printListenError(ListenV1Error error) { + System.err.println("Listen error (" + error.getVariant() + "): " + error.getDescription()); + error.getCode().ifPresent(code -> System.err.println("Error code: " + code)); + } +} diff --git a/src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java b/src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java index 294316bb..b51aacc4 100644 --- a/src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java +++ b/src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java @@ -27,7 +27,6 @@ * Provides production-ready resilience for WebSocket connections. */ public abstract class ReconnectingWebSocketListener extends WebSocketListener { - // A single volatile reference keeps an override internally consistent. private volatile ReconnectOptions activeOptions; private final int maxEnqueuedMessages; @@ -129,7 +128,8 @@ public void connect() { } catch (TimeoutException e) { connectionFuture.cancel(true); TimeoutException timeoutError = - new TimeoutException("WebSocket connection timeout after " + options.connectionTimeoutMs + " milliseconds" + new TimeoutException("WebSocket connection timeout after " + options.connectionTimeoutMs + + " milliseconds" + (retryCount.get() > 0 ? " (retry attempt #" + retryCount.get() + ")" : " (initial connection attempt)")); diff --git a/src/main/java/com/deepgram/resources/agent/v1/types/AgentV1CustomFromThinkProvider.java b/src/main/java/com/deepgram/resources/agent/v1/types/AgentV1CustomFromThinkProvider.java new file mode 100644 index 00000000..777078ce --- /dev/null +++ b/src/main/java/com/deepgram/resources/agent/v1/types/AgentV1CustomFromThinkProvider.java @@ -0,0 +1,125 @@ +/** + * This file was auto-generated by Fern from our API Definition. + */ +package com.deepgram.resources.agent.v1.types; + +import com.deepgram.core.ObjectMappers; +import com.fasterxml.jackson.annotation.JsonAnyGetter; +import com.fasterxml.jackson.annotation.JsonAnySetter; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.annotation.JsonSetter; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import java.util.HashMap; +import java.util.Map; +import java.util.Objects; + +@JsonInclude(JsonInclude.Include.NON_ABSENT) +@JsonDeserialize(builder = AgentV1CustomFromThinkProvider.Builder.class) +public final class AgentV1CustomFromThinkProvider { + private final Object content; + + private final Map additionalProperties; + + private AgentV1CustomFromThinkProvider(Object content, Map additionalProperties) { + this.content = content; + this.additionalProperties = additionalProperties; + } + + /** + * @return Message type identifier for a custom payload returned by the think provider + */ + @JsonProperty("type") + public String getType() { + return "__customFromThinkProvider"; + } + + @JsonProperty("content") + public Object getContent() { + return content; + } + + @java.lang.Override + public boolean equals(Object other) { + if (this == other) return true; + return other instanceof AgentV1CustomFromThinkProvider && equalTo((AgentV1CustomFromThinkProvider) other); + } + + @JsonAnyGetter + public Map getAdditionalProperties() { + return this.additionalProperties; + } + + private boolean equalTo(AgentV1CustomFromThinkProvider other) { + return content.equals(other.content); + } + + @java.lang.Override + public int hashCode() { + return Objects.hash(this.content); + } + + @java.lang.Override + public String toString() { + return ObjectMappers.stringify(this); + } + + public static ContentStage builder() { + return new Builder(); + } + + public interface ContentStage { + _FinalStage content(Object content); + + Builder from(AgentV1CustomFromThinkProvider other); + } + + public interface _FinalStage { + AgentV1CustomFromThinkProvider build(); + + _FinalStage additionalProperty(String key, Object value); + + _FinalStage additionalProperties(Map additionalProperties); + } + + @JsonIgnoreProperties(ignoreUnknown = true) + public static final class Builder implements ContentStage, _FinalStage { + private Object content; + + @JsonAnySetter + private Map additionalProperties = new HashMap<>(); + + private Builder() {} + + @java.lang.Override + public Builder from(AgentV1CustomFromThinkProvider other) { + content(other.getContent()); + return this; + } + + @java.lang.Override + @JsonSetter("content") + public _FinalStage content(Object content) { + this.content = content; + return this; + } + + @java.lang.Override + public AgentV1CustomFromThinkProvider build() { + return new AgentV1CustomFromThinkProvider(content, additionalProperties); + } + + @java.lang.Override + public Builder additionalProperty(String key, Object value) { + this.additionalProperties.put(key, value); + return this; + } + + @java.lang.Override + public Builder additionalProperties(Map additionalProperties) { + this.additionalProperties.putAll(additionalProperties); + return this; + } + } +} diff --git a/src/main/java/com/deepgram/resources/agent/v1/types/AgentV1CustomToThinkProvider.java b/src/main/java/com/deepgram/resources/agent/v1/types/AgentV1CustomToThinkProvider.java new file mode 100644 index 00000000..5dc5b4b1 --- /dev/null +++ b/src/main/java/com/deepgram/resources/agent/v1/types/AgentV1CustomToThinkProvider.java @@ -0,0 +1,125 @@ +/** + * This file was auto-generated by Fern from our API Definition. + */ +package com.deepgram.resources.agent.v1.types; + +import com.deepgram.core.ObjectMappers; +import com.fasterxml.jackson.annotation.JsonAnyGetter; +import com.fasterxml.jackson.annotation.JsonAnySetter; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.annotation.JsonSetter; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import java.util.HashMap; +import java.util.Map; +import java.util.Objects; + +@JsonInclude(JsonInclude.Include.NON_ABSENT) +@JsonDeserialize(builder = AgentV1CustomToThinkProvider.Builder.class) +public final class AgentV1CustomToThinkProvider { + private final Object content; + + private final Map additionalProperties; + + private AgentV1CustomToThinkProvider(Object content, Map additionalProperties) { + this.content = content; + this.additionalProperties = additionalProperties; + } + + /** + * @return Message type identifier for sending a custom payload to the think provider + */ + @JsonProperty("type") + public String getType() { + return "__customToThinkProvider"; + } + + @JsonProperty("content") + public Object getContent() { + return content; + } + + @java.lang.Override + public boolean equals(Object other) { + if (this == other) return true; + return other instanceof AgentV1CustomToThinkProvider && equalTo((AgentV1CustomToThinkProvider) other); + } + + @JsonAnyGetter + public Map getAdditionalProperties() { + return this.additionalProperties; + } + + private boolean equalTo(AgentV1CustomToThinkProvider other) { + return content.equals(other.content); + } + + @java.lang.Override + public int hashCode() { + return Objects.hash(this.content); + } + + @java.lang.Override + public String toString() { + return ObjectMappers.stringify(this); + } + + public static ContentStage builder() { + return new Builder(); + } + + public interface ContentStage { + _FinalStage content(Object content); + + Builder from(AgentV1CustomToThinkProvider other); + } + + public interface _FinalStage { + AgentV1CustomToThinkProvider build(); + + _FinalStage additionalProperty(String key, Object value); + + _FinalStage additionalProperties(Map additionalProperties); + } + + @JsonIgnoreProperties(ignoreUnknown = true) + public static final class Builder implements ContentStage, _FinalStage { + private Object content; + + @JsonAnySetter + private Map additionalProperties = new HashMap<>(); + + private Builder() {} + + @java.lang.Override + public Builder from(AgentV1CustomToThinkProvider other) { + content(other.getContent()); + return this; + } + + @java.lang.Override + @JsonSetter("content") + public _FinalStage content(Object content) { + this.content = content; + return this; + } + + @java.lang.Override + public AgentV1CustomToThinkProvider build() { + return new AgentV1CustomToThinkProvider(content, additionalProperties); + } + + @java.lang.Override + public Builder additionalProperty(String key, Object value) { + this.additionalProperties.put(key, value); + return this; + } + + @java.lang.Override + public Builder additionalProperties(Map additionalProperties) { + this.additionalProperties.putAll(additionalProperties); + return this; + } + } +} diff --git a/src/main/java/com/deepgram/resources/agent/v1/websocket/V1WebSocketClient.java b/src/main/java/com/deepgram/resources/agent/v1/websocket/V1WebSocketClient.java index 5bf4d734..1aa1b9d4 100644 --- a/src/main/java/com/deepgram/resources/agent/v1/websocket/V1WebSocketClient.java +++ b/src/main/java/com/deepgram/resources/agent/v1/websocket/V1WebSocketClient.java @@ -13,6 +13,8 @@ import com.deepgram.resources.agent.v1.types.AgentV1AgentStartedSpeaking; import com.deepgram.resources.agent.v1.types.AgentV1AgentThinking; import com.deepgram.resources.agent.v1.types.AgentV1ConversationText; +import com.deepgram.resources.agent.v1.types.AgentV1CustomFromThinkProvider; +import com.deepgram.resources.agent.v1.types.AgentV1CustomToThinkProvider; import com.deepgram.resources.agent.v1.types.AgentV1Error; import com.deepgram.resources.agent.v1.types.AgentV1ForceEndTurn; import com.deepgram.resources.agent.v1.types.AgentV1FunctionCallCancelled; @@ -111,6 +113,8 @@ public class V1WebSocketClient implements AutoCloseable { private volatile Consumer agentAudioDoneHandler; + private volatile Consumer customFromThinkProviderHandler; + private volatile Consumer errorHandler; private volatile Consumer warningHandler; @@ -334,6 +338,15 @@ public CompletableFuture sendForceEndTurn(AgentV1ForceEndTurn message) { return sendMessage(message); } + /** + * Sends an AgentV1CustomToThinkProvider message to the server asynchronously. + * @param message the message to send + * @return a CompletableFuture that completes when the message is sent + */ + public CompletableFuture sendCustomToThinkProvider(AgentV1CustomToThinkProvider message) { + return sendMessage(message); + } + /** * Sends an AgentV1Media message to the server asynchronously. * @param message the message to send @@ -480,6 +493,14 @@ public void onAgentAudioDone(Consumer handler) { this.agentAudioDoneHandler = handler; } + /** + * Registers a handler for AgentV1CustomFromThinkProvider messages from the server. + * @param handler the handler to invoke when a message is received + */ + public void onCustomFromThinkProvider(Consumer handler) { + this.customFromThinkProviderHandler = handler; + } + /** * Registers a handler for AgentV1Error messages from the server. * @param handler the handler to invoke when a message is received @@ -748,6 +769,21 @@ private void handleIncomingMessage(String json) { return; } } + if (node.has("content") + && "__customFromThinkProvider".equals(node.path("type").asText())) { + AgentV1CustomFromThinkProvider customFromThinkProviderHandlerEvent = null; + try { + customFromThinkProviderHandlerEvent = + objectMapper.treeToValue(node, AgentV1CustomFromThinkProvider.class); + } catch (Exception e) { + } + if (customFromThinkProviderHandlerEvent != null) { + if (customFromThinkProviderHandler != null) { + customFromThinkProviderHandler.accept(customFromThinkProviderHandlerEvent); + } + return; + } + } if ("ListenUpdated".equals(node.path("type").asText())) { AgentV1ListenUpdated listenUpdatedHandlerEvent = null; try { diff --git a/src/main/java/com/deepgram/resources/listen/v1/types/ListenV1Configure.java b/src/main/java/com/deepgram/resources/listen/v1/types/ListenV1Configure.java new file mode 100644 index 00000000..967d27e9 --- /dev/null +++ b/src/main/java/com/deepgram/resources/listen/v1/types/ListenV1Configure.java @@ -0,0 +1,165 @@ +/** + * This file was auto-generated by Fern from our API Definition. + */ +package com.deepgram.resources.listen.v1.types; + +import com.deepgram.core.ObjectMappers; +import com.fasterxml.jackson.annotation.JsonAnyGetter; +import com.fasterxml.jackson.annotation.JsonAnySetter; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.annotation.JsonSetter; +import com.fasterxml.jackson.annotation.Nulls; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; + +@JsonInclude(JsonInclude.Include.NON_ABSENT) +@JsonDeserialize(builder = ListenV1Configure.Builder.class) +public final class ListenV1Configure { + private final Optional> keyterms; + + private final Optional> features; + + private final Map additionalProperties; + + private ListenV1Configure( + Optional> keyterms, + Optional> features, + Map additionalProperties) { + this.keyterms = keyterms; + this.features = features; + this.additionalProperties = additionalProperties; + } + + /** + * @return Message type identifier + */ + @JsonProperty("type") + public String getType() { + return "Configure"; + } + + /** + * @return Replaces the stream's keyterms. Compatible with all Nova-3 models, monolingual and multilingual; on other + * models the server returns an Error with code KeytermsNotSupported. Each array replaces the entire list, + * including keyterms set with the keyterm query parameter. Send an empty array to clear all keyterms. Omit + * the field, or set it to null, to keep the current keyterms. + *

Each entry is a plain term or phrase with no weights or intensifiers. The 500-token keyterm limit that + * applies to the keyterm query parameter also applies to each update. An over-limit update returns an + * Error, and the stream keeps its previous keyterms.

+ */ + @JsonProperty("keyterms") + public Optional> getKeyterms() { + return keyterms; + } + + /** + * @return Turns formatting features on or off. Each key is a feature name and each value is a boolean, for + * example {"numerals": true}. + */ + @JsonProperty("features") + public Optional> getFeatures() { + return features; + } + + @java.lang.Override + public boolean equals(Object other) { + if (this == other) return true; + return other instanceof ListenV1Configure && equalTo((ListenV1Configure) other); + } + + @JsonAnyGetter + public Map getAdditionalProperties() { + return this.additionalProperties; + } + + private boolean equalTo(ListenV1Configure other) { + return keyterms.equals(other.keyterms) && features.equals(other.features); + } + + @java.lang.Override + public int hashCode() { + return Objects.hash(this.keyterms, this.features); + } + + @java.lang.Override + public String toString() { + return ObjectMappers.stringify(this); + } + + public static Builder builder() { + return new Builder(); + } + + @JsonIgnoreProperties(ignoreUnknown = true) + public static final class Builder { + private Optional> keyterms = Optional.empty(); + + private Optional> features = Optional.empty(); + + @JsonAnySetter + private Map additionalProperties = new HashMap<>(); + + private Builder() {} + + public Builder from(ListenV1Configure other) { + keyterms(other.getKeyterms()); + features(other.getFeatures()); + return this; + } + + /** + *

Replaces the stream's keyterms. Compatible with all Nova-3 models, monolingual and multilingual; on other + * models the server returns an Error with code KeytermsNotSupported. Each array replaces the entire list, + * including keyterms set with the keyterm query parameter. Send an empty array to clear all keyterms. Omit + * the field, or set it to null, to keep the current keyterms.

+ *

Each entry is a plain term or phrase with no weights or intensifiers. The 500-token keyterm limit that + * applies to the keyterm query parameter also applies to each update. An over-limit update returns an + * Error, and the stream keeps its previous keyterms.

+ */ + @JsonSetter(value = "keyterms", nulls = Nulls.SKIP) + public Builder keyterms(Optional> keyterms) { + this.keyterms = keyterms; + return this; + } + + public Builder keyterms(List keyterms) { + this.keyterms = Optional.ofNullable(keyterms); + return this; + } + + /** + *

Turns formatting features on or off. Each key is a feature name and each value is a boolean, for + * example {"numerals": true}.

+ */ + @JsonSetter(value = "features", nulls = Nulls.SKIP) + public Builder features(Optional> features) { + this.features = features; + return this; + } + + public Builder features(Map features) { + this.features = Optional.ofNullable(features); + return this; + } + + public ListenV1Configure build() { + return new ListenV1Configure(keyterms, features, additionalProperties); + } + + public Builder additionalProperty(String key, Object value) { + this.additionalProperties.put(key, value); + return this; + } + + public Builder additionalProperties(Map additionalProperties) { + this.additionalProperties.putAll(additionalProperties); + return this; + } + } +} diff --git a/src/main/java/com/deepgram/resources/listen/v1/types/ListenV1Error.java b/src/main/java/com/deepgram/resources/listen/v1/types/ListenV1Error.java new file mode 100644 index 00000000..fd0a4e0c --- /dev/null +++ b/src/main/java/com/deepgram/resources/listen/v1/types/ListenV1Error.java @@ -0,0 +1,267 @@ +/** + * This file was auto-generated by Fern from our API Definition. + */ +package com.deepgram.resources.listen.v1.types; + +import com.deepgram.core.ObjectMappers; +import com.fasterxml.jackson.annotation.JsonAnyGetter; +import com.fasterxml.jackson.annotation.JsonAnySetter; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.annotation.JsonSetter; +import com.fasterxml.jackson.annotation.Nulls; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import java.util.HashMap; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import org.jetbrains.annotations.NotNull; + +@JsonInclude(JsonInclude.Include.NON_ABSENT) +@JsonDeserialize(builder = ListenV1Error.Builder.class) +public final class ListenV1Error { + private final String variant; + + private final String description; + + private final Optional code; + + private final Optional message; + + private final Map additionalProperties; + + private ListenV1Error( + String variant, + String description, + Optional code, + Optional message, + Map additionalProperties) { + this.variant = variant; + this.description = description; + this.code = code; + this.message = message; + this.additionalProperties = additionalProperties; + } + + /** + * @return Message type identifier for error responses + */ + @JsonProperty("type") + public String getType() { + return "Error"; + } + + /** + * @return The error category. SchemaError means the server could not parse a client text message. + * InvalidConfigureMessage means the server rejected a Configure message. + */ + @JsonProperty("variant") + public String getVariant() { + return variant; + } + + /** + * @return A human-readable description of what went wrong + */ + @JsonProperty("description") + public String getDescription() { + return description; + } + + /** + * @return Identifies the reason for an InvalidConfigureMessage error, for example KeytermsNotSupported + * when keyterms are sent on a model other than Nova-3. + */ + @JsonProperty("code") + public Optional getCode() { + return code; + } + + /** + * @return For SchemaError, the original client message that could not be parsed + */ + @JsonProperty("message") + public Optional getMessage() { + return message; + } + + @java.lang.Override + public boolean equals(Object other) { + if (this == other) return true; + return other instanceof ListenV1Error && equalTo((ListenV1Error) other); + } + + @JsonAnyGetter + public Map getAdditionalProperties() { + return this.additionalProperties; + } + + private boolean equalTo(ListenV1Error other) { + return variant.equals(other.variant) + && description.equals(other.description) + && code.equals(other.code) + && message.equals(other.message); + } + + @java.lang.Override + public int hashCode() { + return Objects.hash(this.variant, this.description, this.code, this.message); + } + + @java.lang.Override + public String toString() { + return ObjectMappers.stringify(this); + } + + public static VariantStage builder() { + return new Builder(); + } + + public interface VariantStage { + /** + *

The error category. SchemaError means the server could not parse a client text message. + * InvalidConfigureMessage means the server rejected a Configure message.

+ */ + DescriptionStage variant(@NotNull String variant); + + Builder from(ListenV1Error other); + } + + public interface DescriptionStage { + /** + *

A human-readable description of what went wrong

+ */ + _FinalStage description(@NotNull String description); + } + + public interface _FinalStage { + ListenV1Error build(); + + _FinalStage additionalProperty(String key, Object value); + + _FinalStage additionalProperties(Map additionalProperties); + + /** + *

Identifies the reason for an InvalidConfigureMessage error, for example KeytermsNotSupported + * when keyterms are sent on a model other than Nova-3.

+ */ + _FinalStage code(Optional code); + + _FinalStage code(String code); + + /** + *

For SchemaError, the original client message that could not be parsed

+ */ + _FinalStage message(Optional message); + + _FinalStage message(String message); + } + + @JsonIgnoreProperties(ignoreUnknown = true) + public static final class Builder implements VariantStage, DescriptionStage, _FinalStage { + private String variant; + + private String description; + + private Optional message = Optional.empty(); + + private Optional code = Optional.empty(); + + @JsonAnySetter + private Map additionalProperties = new HashMap<>(); + + private Builder() {} + + @java.lang.Override + public Builder from(ListenV1Error other) { + variant(other.getVariant()); + description(other.getDescription()); + code(other.getCode()); + message(other.getMessage()); + return this; + } + + /** + *

The error category. SchemaError means the server could not parse a client text message. + * InvalidConfigureMessage means the server rejected a Configure message.

+ * @return Reference to {@code this} so that method calls can be chained together. + */ + @java.lang.Override + @JsonSetter("variant") + public DescriptionStage variant(@NotNull String variant) { + this.variant = Objects.requireNonNull(variant, "variant must not be null"); + return this; + } + + /** + *

A human-readable description of what went wrong

+ * @return Reference to {@code this} so that method calls can be chained together. + */ + @java.lang.Override + @JsonSetter("description") + public _FinalStage description(@NotNull String description) { + this.description = Objects.requireNonNull(description, "description must not be null"); + return this; + } + + /** + *

For SchemaError, the original client message that could not be parsed

+ * @return Reference to {@code this} so that method calls can be chained together. + */ + @java.lang.Override + public _FinalStage message(String message) { + this.message = Optional.ofNullable(message); + return this; + } + + /** + *

For SchemaError, the original client message that could not be parsed

+ */ + @java.lang.Override + @JsonSetter(value = "message", nulls = Nulls.SKIP) + public _FinalStage message(Optional message) { + this.message = message; + return this; + } + + /** + *

Identifies the reason for an InvalidConfigureMessage error, for example KeytermsNotSupported + * when keyterms are sent on a model other than Nova-3.

+ * @return Reference to {@code this} so that method calls can be chained together. + */ + @java.lang.Override + public _FinalStage code(String code) { + this.code = Optional.ofNullable(code); + return this; + } + + /** + *

Identifies the reason for an InvalidConfigureMessage error, for example KeytermsNotSupported + * when keyterms are sent on a model other than Nova-3.

+ */ + @java.lang.Override + @JsonSetter(value = "code", nulls = Nulls.SKIP) + public _FinalStage code(Optional code) { + this.code = code; + return this; + } + + @java.lang.Override + public ListenV1Error build() { + return new ListenV1Error(variant, description, code, message, additionalProperties); + } + + @java.lang.Override + public Builder additionalProperty(String key, Object value) { + this.additionalProperties.put(key, value); + return this; + } + + @java.lang.Override + public Builder additionalProperties(Map additionalProperties) { + this.additionalProperties.putAll(additionalProperties); + return this; + } + } +} diff --git a/src/main/java/com/deepgram/resources/listen/v1/websocket/ListenV1ErrorException.java b/src/main/java/com/deepgram/resources/listen/v1/websocket/ListenV1ErrorException.java new file mode 100644 index 00000000..88a946e0 --- /dev/null +++ b/src/main/java/com/deepgram/resources/listen/v1/websocket/ListenV1ErrorException.java @@ -0,0 +1,31 @@ +package com.deepgram.resources.listen.v1.websocket; + +import com.deepgram.core.DeepgramApiException; +import com.deepgram.resources.listen.v1.types.ListenV1Error; + +/** + * Carries a Listen V1 server {@code Error} message to the generic {@code onError} handler when no + * {@code onErrorMessage} handler is registered, so callers on the generic path can still read the + * typed {@link ListenV1Error} fields. + */ +public final class ListenV1ErrorException extends DeepgramApiException { + private final ListenV1Error error; + + public ListenV1ErrorException(ListenV1Error error) { + super(buildMessage(error)); + this.error = error; + } + + /** The server error message as received, with {@code variant}, {@code description}, and optional {@code code}. */ + public ListenV1Error getError() { + return error; + } + + private static String buildMessage(ListenV1Error error) { + String message = error.getVariant() + ": " + error.getDescription(); + if (error.getCode().isPresent()) { + message += " (" + error.getCode().get() + ")"; + } + return message; + } +} diff --git a/src/main/java/com/deepgram/resources/listen/v1/websocket/V1WebSocketClient.java b/src/main/java/com/deepgram/resources/listen/v1/websocket/V1WebSocketClient.java index 278a2c9b..e7cd50c4 100644 --- a/src/main/java/com/deepgram/resources/listen/v1/websocket/V1WebSocketClient.java +++ b/src/main/java/com/deepgram/resources/listen/v1/websocket/V1WebSocketClient.java @@ -11,6 +11,8 @@ import com.deepgram.core.RequestOptions; import com.deepgram.core.WebSocketReadyState; import com.deepgram.resources.listen.v1.types.ListenV1CloseStream; +import com.deepgram.resources.listen.v1.types.ListenV1Configure; +import com.deepgram.resources.listen.v1.types.ListenV1Error; import com.deepgram.resources.listen.v1.types.ListenV1Finalize; import com.deepgram.resources.listen.v1.types.ListenV1KeepAlive; import com.deepgram.resources.listen.v1.types.ListenV1Metadata; @@ -66,6 +68,8 @@ public class V1WebSocketClient implements AutoCloseable { private volatile Consumer speechStartedHandler; + private volatile Consumer errorHandler; + /** * Creates a new async WebSocket client for the v1 channel. */ @@ -352,6 +356,15 @@ public CompletableFuture sendKeepAlive(ListenV1KeepAlive message) { return sendMessage(message); } + /** + * Sends a ListenV1Configure message to the server asynchronously. + * @param message the message to send + * @return a CompletableFuture that completes when the message is sent + */ + public CompletableFuture sendConfigure(ListenV1Configure message) { + return sendMessage(message); + } + /** * Registers a handler for ListenV1Results messages from the server. * @param handler the handler to invoke when a message is received @@ -384,6 +397,14 @@ public void onSpeechStarted(Consumer handler) { this.speechStartedHandler = handler; } + /** + * Registers a handler for ListenV1Error messages from the server. + * @param handler the handler to invoke when a message is received + */ + public void onErrorMessage(Consumer handler) { + this.errorHandler = handler; + } + /** * Registers a handler called when the connection is established. * @param handler the handler to invoke when connected @@ -540,6 +561,23 @@ private void handleIncomingMessage(String json) { return; } } + if (node.has("variant") + && node.has("description") + && "Error".equals(node.path("type").asText())) { + ListenV1Error errorHandlerEvent = null; + try { + errorHandlerEvent = objectMapper.treeToValue(node, ListenV1Error.class); + } catch (Exception e) { + } + if (errorHandlerEvent != null) { + if (errorHandler != null) { + errorHandler.accept(errorHandlerEvent); + } else if (onErrorHandler != null) { + onErrorHandler.accept(new ListenV1ErrorException(errorHandlerEvent)); + } + return; + } + } if (onErrorHandler != null) { onErrorHandler.accept(new RuntimeException( "Unrecognized WebSocket message: " + json.substring(0, Math.min(200, json.length())) diff --git a/src/test/java/com/deepgram/AgentV1ControlFrameWireTest.java b/src/test/java/com/deepgram/AgentV1ControlFrameWireTest.java index 4404a345..1e5b7f3b 100644 --- a/src/test/java/com/deepgram/AgentV1ControlFrameWireTest.java +++ b/src/test/java/com/deepgram/AgentV1ControlFrameWireTest.java @@ -3,13 +3,20 @@ import static org.assertj.core.api.Assertions.assertThat; import com.deepgram.core.Environment; +import com.deepgram.core.ObjectMappers; +import com.deepgram.resources.agent.v1.types.AgentV1CustomFromThinkProvider; +import com.deepgram.resources.agent.v1.types.AgentV1CustomToThinkProvider; import com.deepgram.resources.agent.v1.types.AgentV1FunctionCallCancelled; import com.deepgram.resources.agent.v1.types.AgentV1ForceEndTurn; import com.deepgram.resources.agent.v1.websocket.V1WebSocketClient; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; import java.util.concurrent.BlockingQueue; import java.util.concurrent.CountDownLatch; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import okhttp3.WebSocket; import okhttp3.WebSocketListener; @@ -97,4 +104,105 @@ public void onClosing(WebSocket webSocket, int code, String reason) { ws.disconnect(); } } + + @Test + void sendsCustomPayloadToThinkProvider() throws Exception { + BlockingQueue received = new LinkedBlockingQueue<>(); + server.enqueue(new MockResponse().withWebSocketUpgrade(new WebSocketListener() { + @Override + public void onMessage(WebSocket webSocket, String text) { + received.add(text); + } + + @Override + public void onClosing(WebSocket webSocket, int code, String reason) { + webSocket.close(code, reason); + } + })); + + V1WebSocketClient ws = client.agent().v1().v1WebSocket(); + try { + ws.connect().get(5, TimeUnit.SECONDS); + Map content = new LinkedHashMap<>(); + content.put("action", "lookup"); + content.put("arguments", Map.of("city", "Seattle", "units", "metric")); + content.put("candidates", List.of("weather", 7, true)); + content.put("attempt", 2); + content.put("enabled", false); + content.put("empty", null); + ws.sendCustomToThinkProvider(AgentV1CustomToThinkProvider.builder() + .content(content) + .build()) + .get(5, TimeUnit.SECONDS); + + String frame = received.poll(5, TimeUnit.SECONDS); + assertThat(frame).isNotNull(); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("type").asText()) + .isEqualTo("__customToThinkProvider"); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("content").path("action").asText()) + .isEqualTo("lookup"); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("content").path("arguments").path("city").asText()) + .isEqualTo("Seattle"); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("content").path("candidates").get(1).asInt()) + .isEqualTo(7); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("content").path("candidates").get(2).asBoolean()) + .isTrue(); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("content").path("attempt").asInt()) + .isEqualTo(2); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("content").path("enabled").asBoolean()) + .isFalse(); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("content").path("empty").isNull()) + .isTrue(); + } finally { + ws.disconnect(); + } + } + + @Test + void dispatchesCustomPayloadFromThinkProvider() throws Exception { + CountDownLatch received = new CountDownLatch(1); + AtomicReference custom = new AtomicReference<>(); + AtomicInteger genericErrorCount = new AtomicInteger(); + server.enqueue(new MockResponse().withWebSocketUpgrade(new WebSocketListener() { + @Override + public void onOpen(WebSocket webSocket, okhttp3.Response response) { + webSocket.send("{\"type\":\"__customFromThinkProvider\",\"content\":{" + + "\"decision\":\"continue\",\"arguments\":{\"city\":\"Seattle\"}," + + "\"candidates\":[\"weather\",7,true],\"attempt\":2,\"enabled\":false,\"empty\":null}}"); + } + + @Override + public void onClosing(WebSocket webSocket, int code, String reason) { + webSocket.close(code, reason); + } + })); + + V1WebSocketClient ws = client.agent().v1().v1WebSocket(); + ws.onCustomFromThinkProvider(event -> { + custom.set(event); + received.countDown(); + }); + ws.onError(event -> genericErrorCount.incrementAndGet()); + try { + ws.connect().get(5, TimeUnit.SECONDS); + + assertThat(received.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(custom.get().getContent()).isInstanceOf(Map.class); + Map content = (Map) custom.get().getContent(); + assertThat(content.get("decision")).isEqualTo("continue"); + Map arguments = (Map) content.get("arguments"); + assertThat(arguments.get("city")).isEqualTo("Seattle"); + List candidates = (List) content.get("candidates"); + assertThat(candidates).hasSize(3); + assertThat(candidates.get(0)).isEqualTo("weather"); + assertThat(candidates.get(1)).isEqualTo(7); + assertThat(candidates.get(2)).isEqualTo(true); + assertThat(content.get("attempt")).isEqualTo(2); + assertThat(content.get("enabled")).isEqualTo(false); + assertThat(content.get("empty")).isNull(); + assertThat(genericErrorCount).hasValue(0); + } finally { + ws.disconnect(); + } + } } diff --git a/src/test/java/com/deepgram/ListenV1ConfigureErrorWebSocketTest.java b/src/test/java/com/deepgram/ListenV1ConfigureErrorWebSocketTest.java new file mode 100644 index 00000000..3e89b974 --- /dev/null +++ b/src/test/java/com/deepgram/ListenV1ConfigureErrorWebSocketTest.java @@ -0,0 +1,199 @@ +package com.deepgram; + +import static org.assertj.core.api.Assertions.assertThat; + +import com.deepgram.core.Environment; +import com.deepgram.core.ObjectMappers; +import com.deepgram.resources.listen.v1.types.ListenV1Configure; +import com.deepgram.resources.listen.v1.types.ListenV1Error; +import com.deepgram.resources.listen.v1.websocket.V1ConnectOptions; +import com.deepgram.resources.listen.v1.websocket.ListenV1ErrorException; +import com.deepgram.resources.listen.v1.websocket.V1WebSocketClient; +import com.deepgram.types.ListenV1Model; +import java.util.List; +import java.util.Map; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; +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.Test; + +class ListenV1ConfigureErrorWebSocketTest { + 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 + void sendsConfigureAndDispatchesConfigureError() throws Exception { + BlockingQueue received = new LinkedBlockingQueue<>(); + server.enqueue(new MockResponse().withWebSocketUpgrade(new WebSocketListener() { + @Override + public void onMessage(WebSocket webSocket, String text) { + received.add(text); + webSocket.send("{\"type\":\"Error\",\"variant\":\"InvalidConfigureMessage\"," + + "\"description\":\"keyterms unavailable\",\"code\":\"KeytermsNotSupported\"}"); + } + + @Override + public void onClosing(WebSocket webSocket, int code, String reason) { + webSocket.close(code, reason); + } + })); + + CountDownLatch errorReceived = new CountDownLatch(1); + AtomicInteger genericErrorCount = new AtomicInteger(); + AtomicReference error = new AtomicReference<>(); + V1WebSocketClient ws = client.listen().v1().v1WebSocket(); + ws.onErrorMessage(event -> { + error.set(event); + errorReceived.countDown(); + }); + ws.onError(event -> genericErrorCount.incrementAndGet()); + try { + ws.connect(V1ConnectOptions.builder().model(ListenV1Model.NOVA3).build()) + .get(5, TimeUnit.SECONDS); + ws.sendConfigure(ListenV1Configure.builder() + .keyterms(List.of("Deepgram", "Flux")) + .features(Map.of("numerals", false)) + .build()) + .get(5, TimeUnit.SECONDS); + + String frame = received.poll(5, TimeUnit.SECONDS); + assertThat(frame).isNotNull(); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("type").asText()).isEqualTo("Configure"); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("keyterms").get(0).asText()) + .isEqualTo("Deepgram"); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("keyterms").get(1).asText()) + .isEqualTo("Flux"); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("features").path("numerals").asBoolean()) + .isFalse(); + + assertThat(errorReceived.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(error.get().getVariant()).isEqualTo("InvalidConfigureMessage"); + assertThat(error.get().getCode()).contains("KeytermsNotSupported"); + assertThat(genericErrorCount).hasValue(0); + + ws.sendConfigure(ListenV1Configure.builder().keyterms(List.of()).build()) + .get(5, TimeUnit.SECONDS); + frame = received.poll(5, TimeUnit.SECONDS); + assertThat(frame).isNotNull(); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("keyterms").isArray()).isTrue(); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("keyterms").isEmpty()).isTrue(); + + ws.sendConfigure(ListenV1Configure.builder().build()).get(5, TimeUnit.SECONDS); + frame = received.poll(5, TimeUnit.SECONDS); + assertThat(frame).isNotNull(); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).size()).isEqualTo(1); + assertThat(ObjectMappers.JSON_MAPPER.readTree(frame).path("type").asText()).isEqualTo("Configure"); + } finally { + ws.disconnect(); + } + } + + @Test + void dispatchesSchemaErrorWithMessageAndNoCode() throws Exception { + server.enqueue(new MockResponse().withWebSocketUpgrade(new WebSocketListener() { + @Override + public void onMessage(WebSocket webSocket, String text) { + webSocket.send("{\"type\":\"Error\",\"variant\":\"SchemaError\"," + + "\"description\":\"could not parse control frame\"," + + "\"message\":\"{\\\"type\\\":\\\"Configure\\\"}\"}"); + } + + @Override + public void onClosing(WebSocket webSocket, int code, String reason) { + webSocket.close(code, reason); + } + })); + + CountDownLatch errorReceived = new CountDownLatch(1); + AtomicInteger genericErrorCount = new AtomicInteger(); + AtomicReference error = new AtomicReference<>(); + V1WebSocketClient ws = client.listen().v1().v1WebSocket(); + ws.onErrorMessage(event -> { + error.set(event); + errorReceived.countDown(); + }); + ws.onError(event -> genericErrorCount.incrementAndGet()); + try { + ws.connect(V1ConnectOptions.builder().model(ListenV1Model.NOVA3).build()) + .get(5, TimeUnit.SECONDS); + ws.sendConfigure(ListenV1Configure.builder().build()).get(5, TimeUnit.SECONDS); + + assertThat(errorReceived.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(error.get().getVariant()).isEqualTo("SchemaError"); + assertThat(error.get().getDescription()).isEqualTo("could not parse control frame"); + assertThat(error.get().getMessage()).contains("{\"type\":\"Configure\"}"); + assertThat(error.get().getCode()).isEmpty(); + assertThat(genericErrorCount).hasValue(0); + } finally { + ws.disconnect(); + } + } + + @Test + void dispatchesServerErrorsToGenericHandlerWithoutTypedHandler() throws Exception { + server.enqueue(new MockResponse().withWebSocketUpgrade(new WebSocketListener() { + @Override + public void onMessage(WebSocket webSocket, String text) { + webSocket.send("{\"type\":\"Error\",\"variant\":\"SchemaError\"," + + "\"description\":\"could not parse control frame\"}"); + } + + @Override + public void onClosing(WebSocket webSocket, int code, String reason) { + webSocket.close(code, reason); + } + })); + + CountDownLatch errorReceived = new CountDownLatch(1); + AtomicReference error = new AtomicReference<>(); + V1WebSocketClient ws = client.listen().v1().v1WebSocket(); + ws.onError(event -> { + error.set(event); + errorReceived.countDown(); + }); + try { + ws.connect(V1ConnectOptions.builder().model(ListenV1Model.NOVA3).build()) + .get(5, TimeUnit.SECONDS); + ws.sendConfigure(ListenV1Configure.builder().build()).get(5, TimeUnit.SECONDS); + + assertThat(errorReceived.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(error.get()) + .isInstanceOf(ListenV1ErrorException.class) + .hasMessage("SchemaError: could not parse control frame"); + ListenV1Error serverError = ((ListenV1ErrorException) error.get()).getError(); + assertThat(serverError.getVariant()).isEqualTo("SchemaError"); + assertThat(serverError.getDescription()).isEqualTo("could not parse control frame"); + assertThat(serverError.getCode()).isEmpty(); + } finally { + ws.disconnect(); + } + } +}