diff --git a/.github/workflows/post-release.yml b/.github/workflows/post-release.yml index 0b1eeb12..f03a1a0a 100644 --- a/.github/workflows/post-release.yml +++ b/.github/workflows/post-release.yml @@ -30,7 +30,7 @@ jobs: - name: Revert pom.xml to ${{env.NEXT_VERSION}} run: | ./mvnw versions:set -DnewVersion=${{env.NEXT_VERSION}} - ./mvnw -f samples/pom.xml versions:use-dep-version -Dincludes=com.ibm.watsonx:watsonx-ai -DdepVersion=${{env.CURRENT_VERSION}} -DforceVersion=true + ./mvnw -f samples/pom.xml versions:use-dep-version -Dincludes=com.ibm.watsonx:* -DdepVersion=${{env.CURRENT_VERSION}} -DforceVersion=true - name: Commit revert uses: ./.github/actions/commit-push diff --git a/.github/workflows/pre-release.yml b/.github/workflows/pre-release.yml index 9e4874ad..6bb02938 100644 --- a/.github/workflows/pre-release.yml +++ b/.github/workflows/pre-release.yml @@ -40,6 +40,7 @@ jobs: run: | ./mvnw versions:set -DnewVersion=${{env.CURRENT_VERSION}} ./mvnw versions:set -f samples/pom.xml -DnewVersion=${{env.CURRENT_VERSION}} + ./mvnw -f samples/pom.xml versions:use-dep-version -Dincludes=com.ibm.watsonx:* -DdepVersion=${{env.CURRENT_VERSION}} -DforceVersion=true - name: Update SDK version in docusaurus.config.ts run: | diff --git a/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/chat/ChatService.java b/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/chat/ChatService.java index ad327815..aab22e9e 100644 --- a/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/chat/ChatService.java +++ b/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/chat/ChatService.java @@ -167,6 +167,44 @@ public CompletableFuture chatStreaming(ChatRequest chatRequest, Ch return client.chatStreaming(transactionId, textChatRequest, context, handler); } + /** + * Sends a streaming chat request. + * + * @param chatRequest the {@link ChatRequest} + * @param onResponse a consumer that receives partial response tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} + */ + public CompletableFuture chatStreaming(ChatRequest chatRequest, Consumer onResponse) { + return chatStreaming(chatRequest, new ChatHandler() { + @Override + public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { + onResponse.accept(partialResponse); + } + }); + } + + /** + * Sends a streaming chat request, delivering response and reasoning tokens to separate consumers. + * + * @param chatRequest the {@link ChatRequest} + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} + */ + public CompletableFuture chatStreaming(ChatRequest chatRequest, Consumer onResponse, Consumer onThinking) { + return chatStreaming(chatRequest, new ChatHandler() { + @Override + public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { + onResponse.accept(partialResponse); + } + + @Override + public void onPartialThinking(String partialThinking, PartialChatResponse partialChatResponse) { + onThinking.accept(partialThinking); + } + }); + } + /** * Sends a chat request to the model using the provided message. * @@ -265,7 +303,7 @@ public CompletableFuture chatStreaming(String message, ChatHandler * @param handler a {@link ChatHandler} implementation */ public CompletableFuture chatStreaming(List messages, ChatHandler handler) { - return chatStreaming(messages, ChatParameters.builder().build(), handler); + return chatStreaming(messages, (ChatParameters) null, handler); } /** @@ -276,7 +314,7 @@ public CompletableFuture chatStreaming(List messages, * @param handler a {@link ChatHandler} implementation */ public CompletableFuture chatStreaming(List messages, List tools, ChatHandler handler) { - return chatStreaming(messages, null, tools, handler); + return chatStreaming(messages, (ChatParameters) null, tools, handler); } /** @@ -287,27 +325,59 @@ public CompletableFuture chatStreaming(List messages, * @param handler a {@link ChatHandler} implementation */ public CompletableFuture chatStreaming(List messages, ChatParameters parameters, ChatHandler handler) { - return chatStreaming(messages, parameters, null, handler); + return chatStreaming(messages, parameters, (List) null, handler); } /** - * Sends a chat request to the model using the provided message. + * Sends a streaming chat request using the provided message. * - * @param message Message to send. - * @param handler a consumer that receives partial text responses + * @param message the message to send + * @param onResponse a consumer that receives partial response tokens */ - public CompletableFuture chatStreaming(String message, Consumer handler) { - return chatStreaming(List.of(UserMessage.text(message)), handler); + public CompletableFuture chatStreaming(String message, Consumer onResponse) { + return chatStreaming(List.of(UserMessage.text(message)), onResponse); } /** * Sends a streaming chat request using the provided messages. * * @param messages the list of chat messages forming the prompt history - * @param handler a consumer that receives partial text responses + * @param onResponse a consumer that receives partial response tokens + */ + public CompletableFuture chatStreaming(List messages, Consumer onResponse) { + return chatStreaming(messages, (ChatParameters) null, onResponse); + } + + /** + * Sends a streaming chat request using the provided message, delivering response and reasoning tokens to separate consumers. + * + * @param message the message to send + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + */ + public CompletableFuture chatStreaming(String message, Consumer onResponse, Consumer onThinking) { + return chatStreaming(List.of(UserMessage.text(message)), onResponse, onThinking); + } + + /** + * Sends a streaming chat request using the provided messages, delivering response and reasoning tokens to separate consumers. + * + * @param messages the list of chat messages forming the prompt history + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens */ - public CompletableFuture chatStreaming(List messages, Consumer handler) { - return chatStreaming(messages, ChatParameters.builder().build(), handler); + public CompletableFuture chatStreaming(List messages, Consumer onResponse, Consumer onThinking) { + return chatStreaming(messages, (ChatParameters) null, (List) null, new ChatHandler() { + @Override + public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { + onResponse.accept(partialResponse); + } + + @Override + public void onPartialThinking(String partialThinking, PartialChatResponse partialChatResponse) { + onThinking.accept(partialThinking); + } + }); } /** @@ -318,7 +388,7 @@ public CompletableFuture chatStreaming(List messages, * @param handler a consumer that receives partial text responses */ public CompletableFuture chatStreaming(List messages, ChatParameters parameters, Consumer handler) { - return chatStreaming(messages, parameters, null, handler); + return chatStreaming(messages, parameters, (List) null, handler); } /** @@ -329,7 +399,7 @@ public CompletableFuture chatStreaming(List messages, * @param handler a consumer that receives partial text responses */ public CompletableFuture chatStreaming(List messages, List tools, Consumer handler) { - return chatStreaming(messages, null, tools, handler); + return chatStreaming(messages, (ChatParameters) null, tools, handler); } /** diff --git a/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/deployment/DeploymentService.java b/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/deployment/DeploymentService.java index c36b591f..1334bb81 100644 --- a/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/deployment/DeploymentService.java +++ b/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/deployment/DeploymentService.java @@ -209,17 +209,39 @@ public TextChatResponse chat(DeploymentChatRequest chatRequest) { /** * Sends a streaming chat request. - *

- * This method initiates an asynchronous chat operation where partial responses are delivered incrementally through the provided {@link Consumer}. * - * @param chatRequest the chat request - * @param handler a consumer that receives partial text responses + * @param chatRequest the {@link DeploymentChatRequest} + * @param onResponse a consumer that receives partial response tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} */ - public CompletableFuture chatStreaming(DeploymentChatRequest chatRequest, Consumer handler) { + public CompletableFuture chatStreaming(DeploymentChatRequest chatRequest, Consumer onResponse) { return chatStreaming(chatRequest, new ChatHandler() { @Override public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { - handler.accept(partialResponse); + onResponse.accept(partialResponse); + } + }); + } + + /** + * Sends a streaming chat request, delivering response and reasoning tokens to separate consumers. + * + * @param chatRequest the {@link DeploymentChatRequest} + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} + */ + public CompletableFuture chatStreaming(DeploymentChatRequest chatRequest, Consumer onResponse, + Consumer onThinking) { + return chatStreaming(chatRequest, new ChatHandler() { + @Override + public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { + onResponse.accept(partialResponse); + } + + @Override + public void onPartialThinking(String partialThinking, PartialChatResponse partialChatResponse) { + onThinking.accept(partialThinking); } }); } @@ -365,7 +387,7 @@ public CompletableFuture chatStreaming(String deploymentId, String * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} */ public CompletableFuture chatStreaming(String deploymentId, List messages, ChatHandler handler) { - return chatStreaming(deploymentId, messages, ChatParameters.builder().build(), handler); + return chatStreaming(deploymentId, messages, (ChatParameters) null, handler); } /** @@ -378,7 +400,7 @@ public CompletableFuture chatStreaming(String deploymentId, List chatStreaming(String deploymentId, List messages, List tools, ChatHandler handler) { - return chatStreaming(deploymentId, messages, null, tools, handler); + return chatStreaming(deploymentId, messages, (ChatParameters) null, tools, handler); } /** @@ -392,19 +414,19 @@ public CompletableFuture chatStreaming(String deploymentId, List chatStreaming(String deploymentId, List messages, ChatParameters parameters, ChatHandler handler) { - return chatStreaming(deploymentId, messages, parameters, null, handler); + return chatStreaming(deploymentId, messages, parameters, (List) null, handler); } /** * Sends a streaming chat request to a deployment using the provided message. * * @param deploymentId the unique identifier of the deployment - * @param message Message to send. - * @param handler a consumer that receives partial text responses + * @param message the message to send + * @param onResponse a consumer that receives partial response tokens * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} */ - public CompletableFuture chatStreaming(String deploymentId, String message, Consumer handler) { - return chatStreaming(deploymentId, List.of(UserMessage.text(message)), handler); + public CompletableFuture chatStreaming(String deploymentId, String message, Consumer onResponse) { + return chatStreaming(deploymentId, List.of(UserMessage.text(message)), onResponse); } /** @@ -412,11 +434,11 @@ public CompletableFuture chatStreaming(String deploymentId, String * * @param deploymentId the unique identifier of the deployment * @param messages the list of chat messages forming the prompt history - * @param handler a consumer that receives partial text responses + * @param onResponse a consumer that receives partial response tokens * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} */ - public CompletableFuture chatStreaming(String deploymentId, List messages, Consumer handler) { - return chatStreaming(deploymentId, messages, ChatParameters.builder().build(), handler); + public CompletableFuture chatStreaming(String deploymentId, List messages, Consumer onResponse) { + return chatStreaming(deploymentId, messages, (ChatParameters) null, onResponse); } /** @@ -444,7 +466,7 @@ public CompletableFuture chatStreaming(String deploymentId, List chatStreaming(String deploymentId, List messages, ChatParameters parameters, Consumer handler) { - return chatStreaming(deploymentId, messages, parameters, null, handler); + return chatStreaming(deploymentId, messages, parameters, (List) null, handler); } /** @@ -467,6 +489,103 @@ public void onPartialResponse(String partialResponse, PartialChatResponse partia }); } + /** + * Sends a streaming chat request to a deployment using the provided message, delivering response and reasoning tokens to separate consumers. + * + * @param deploymentId the unique identifier of the deployment + * @param message the message to send + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} + */ + public CompletableFuture chatStreaming(String deploymentId, String message, Consumer onResponse, + Consumer onThinking) { + return chatStreaming(deploymentId, List.of(UserMessage.text(message)), onResponse, onThinking); + } + + /** + * Sends a streaming chat request to a deployment using the provided messages, delivering response and reasoning tokens to separate consumers. + * + * @param deploymentId the unique identifier of the deployment + * @param messages the list of chat messages forming the prompt history + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} + */ + public CompletableFuture chatStreaming(String deploymentId, List messages, Consumer onResponse, + Consumer onThinking) { + return chatStreaming(deploymentId, messages, (ChatParameters) null, (List) null, new ChatHandler() { + @Override + public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { + onResponse.accept(partialResponse); + } + + @Override + public void onPartialThinking(String partialThinking, PartialChatResponse partialChatResponse) { + onThinking.accept(partialThinking); + } + }); + } + + /** + * Sends a streaming chat request to a deployment using the provided messages and parameters, delivering response and reasoning tokens to separate + * consumers. + * + * @param deploymentId the unique identifier of the deployment + * @param messages the list of chat messages forming the prompt history + * @param parameters additional optional parameters for the chat invocation + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} + */ + public CompletableFuture chatStreaming(String deploymentId, List messages, ChatParameters parameters, + Consumer onResponse, Consumer onThinking) { + return chatStreaming(deploymentId, messages, parameters, (List) null, onResponse, onThinking); + } + + /** + * Sends a streaming chat request to a deployment using the provided messages and tools, delivering response and reasoning tokens to separate + * consumers. + * + * @param deploymentId the unique identifier of the deployment + * @param messages the list of chat messages forming the prompt history + * @param tools the list of tools that the model may use + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} + */ + public CompletableFuture chatStreaming(String deploymentId, List messages, List tools, + Consumer onResponse, Consumer onThinking) { + return chatStreaming(deploymentId, messages, (ChatParameters) null, tools, onResponse, onThinking); + } + + /** + * Sends a streaming chat request to a deployment using the provided messages, parameters and tools, delivering response and reasoning tokens to + * separate consumers. + * + * @param deploymentId the unique identifier of the deployment + * @param messages the list of chat messages forming the prompt history + * @param parameters additional optional parameters for the chat invocation + * @param tools the list of tools that the model may use + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} + */ + public CompletableFuture chatStreaming(String deploymentId, List messages, ChatParameters parameters, + List tools, Consumer onResponse, Consumer onThinking) { + return chatStreaming(deploymentId, messages, parameters, tools, new ChatHandler() { + @Override + public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { + onResponse.accept(partialResponse); + } + + @Override + public void onPartialThinking(String partialThinking, PartialChatResponse partialChatResponse) { + onThinking.accept(partialThinking); + } + }); + } + /** * Sends a streaming chat request to a deployment using the provided messages, parameters and tools. * diff --git a/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/gateway/chat/ModelGatewayChatService.java b/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/gateway/chat/ModelGatewayChatService.java index 933a442c..1a3c0768 100644 --- a/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/gateway/chat/ModelGatewayChatService.java +++ b/modules/watsonx-ai/src/main/java/com/ibm/watsonx/ai/gateway/chat/ModelGatewayChatService.java @@ -150,6 +150,45 @@ public CompletableFuture chatStreaming(ModelGatewayChatRequest cha return client.chatStreaming(transactionId, Duration.ofMillis(gatewayRequest.timeLimit()), gatewayRequest, context, handler); } + /** + * Sends a streaming chat request to the Model Gateway. + * + * @param chatRequest the {@link ModelGatewayChatRequest} + * @param onResponse a consumer that receives partial response tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} + */ + public CompletableFuture chatStreaming(ModelGatewayChatRequest chatRequest, Consumer onResponse) { + return chatStreaming(chatRequest, new ChatHandler() { + @Override + public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { + onResponse.accept(partialResponse); + } + }); + } + + /** + * Sends a streaming chat request to the Model Gateway, delivering response and reasoning tokens to separate consumers. + * + * @param chatRequest the {@link ModelGatewayChatRequest} + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes with the final {@link ChatResponse} + */ + public CompletableFuture chatStreaming(ModelGatewayChatRequest chatRequest, Consumer onResponse, + Consumer onThinking) { + return chatStreaming(chatRequest, new ChatHandler() { + @Override + public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { + onResponse.accept(partialResponse); + } + + @Override + public void onPartialThinking(String partialThinking, PartialChatResponse partialChatResponse) { + onThinking.accept(partialThinking); + } + }); + } + /** * Sends a chat request to the model using the provided message. * @@ -262,7 +301,7 @@ public CompletableFuture chatStreaming(List messages, * @return a {@link CompletableFuture} that completes when the stream finishes or fails */ public CompletableFuture chatStreaming(List messages, List tools, ChatHandler handler) { - return chatStreaming(messages, null, tools, handler); + return chatStreaming(messages, (ModelGatewayChatParameters) null, tools, handler); } /** @@ -274,29 +313,29 @@ public CompletableFuture chatStreaming(List messages, * @return a {@link CompletableFuture} that completes when the stream finishes or fails */ public CompletableFuture chatStreaming(List messages, ModelGatewayChatParameters parameters, ChatHandler handler) { - return chatStreaming(messages, parameters, null, handler); + return chatStreaming(messages, parameters, (List) null, handler); } /** - * Sends a streaming chat request using the provided message, delegating to a simple text consumer. + * Sends a streaming chat request using the provided message. * * @param message the message to send - * @param handler a consumer that receives partial text responses + * @param onResponse a consumer that receives partial response tokens * @return a {@link CompletableFuture} that completes when the stream finishes or fails */ - public CompletableFuture chatStreaming(String message, Consumer handler) { - return chatStreaming(List.of(UserMessage.text(message)), handler); + public CompletableFuture chatStreaming(String message, Consumer onResponse) { + return chatStreaming(List.of(UserMessage.text(message)), onResponse); } /** - * Sends a streaming chat request using the provided messages, delegating to a simple text consumer. + * Sends a streaming chat request using the provided messages. * * @param messages the list of chat messages forming the prompt history - * @param handler a consumer that receives partial text responses + * @param onResponse a consumer that receives partial response tokens * @return a {@link CompletableFuture} that completes when the stream finishes or fails */ - public CompletableFuture chatStreaming(List messages, Consumer handler) { - return chatStreaming(messages, (ModelGatewayChatParameters) null, handler); + public CompletableFuture chatStreaming(List messages, Consumer onResponse) { + return chatStreaming(messages, (ModelGatewayChatParameters) null, onResponse); } /** @@ -309,7 +348,7 @@ public CompletableFuture chatStreaming(List messages, */ public CompletableFuture chatStreaming(List messages, ModelGatewayChatParameters parameters, Consumer handler) { - return chatStreaming(messages, parameters, null, handler); + return chatStreaming(messages, parameters, (List) null, handler); } /** @@ -343,6 +382,96 @@ public void onPartialResponse(String partialResponse, PartialChatResponse partia }); } + /** + * Sends a streaming chat request using the provided message, delivering response and reasoning tokens to separate consumers. + * + * @param message the message to send + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes when the stream finishes or fails + */ + public CompletableFuture chatStreaming(String message, Consumer onResponse, Consumer onThinking) { + return chatStreaming(List.of(UserMessage.text(message)), onResponse, onThinking); + } + + /** + * Sends a streaming chat request using the provided messages, delivering response and reasoning tokens to separate consumers. + * + * @param messages the list of chat messages forming the prompt history + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes when the stream finishes or fails + */ + public CompletableFuture chatStreaming(List messages, Consumer onResponse, Consumer onThinking) { + return chatStreaming(messages, (ModelGatewayChatParameters) null, (List) null, new ChatHandler() { + @Override + public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { + onResponse.accept(partialResponse); + } + + @Override + public void onPartialThinking(String partialThinking, PartialChatResponse partialChatResponse) { + onThinking.accept(partialThinking); + } + }); + } + + /** + * Sends a streaming chat request using the provided messages and gateway parameters, delivering response and reasoning tokens to separate + * consumers. + * + * @param messages the list of chat messages forming the prompt history + * @param parameters gateway parameters for the chat invocation + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes when the stream finishes or fails + */ + public CompletableFuture chatStreaming( + List messages, ModelGatewayChatParameters parameters, Consumer onResponse, Consumer onThinking) { + return chatStreaming(messages, parameters, (List) null, onResponse, onThinking); + } + + /** + * Sends a streaming chat request using the provided messages and tools, delivering response and reasoning tokens to separate consumers. + * + * @param messages the list of chat messages forming the prompt history + * @param tools the list of tools that the model may use + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes when the stream finishes or fails + */ + public CompletableFuture chatStreaming( + List messages, List tools, Consumer onResponse, Consumer onThinking) { + return chatStreaming(messages, (ModelGatewayChatParameters) null, tools, onResponse, onThinking); + } + + /** + * Sends a streaming chat request using the provided messages, parameters, and tools, delivering response and reasoning tokens to separate + * consumers. + * + * @param messages the list of chat messages forming the prompt history + * @param parameters gateway parameters for the chat invocation + * @param tools the list of tools that the model may use + * @param onResponse a consumer that receives partial response tokens + * @param onThinking a consumer that receives partial reasoning tokens + * @return a {@link CompletableFuture} that completes when the stream finishes or fails + */ + public CompletableFuture chatStreaming( + List messages, ModelGatewayChatParameters parameters, List tools, Consumer onResponse, + Consumer onThinking) { + return chatStreaming(messages, parameters, tools, new ChatHandler() { + @Override + public void onPartialResponse(String partialResponse, PartialChatResponse partialChatResponse) { + onResponse.accept(partialResponse); + } + + @Override + public void onPartialThinking(String partialThinking, PartialChatResponse partialChatResponse) { + onThinking.accept(partialThinking); + } + }); + } + /** * Sends a streaming chat request using the provided messages, parameters, and tools. * diff --git a/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/DeploymentServiceTest.java b/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/DeploymentServiceTest.java index b7f81649..822f55ec 100644 --- a/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/DeploymentServiceTest.java +++ b/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/DeploymentServiceTest.java @@ -2540,4 +2540,69 @@ private void assertStreamedContent(CompletableFuture future, Strin assertEquals("Ciao", partial.toString()); partial.setLength(0); } + + @Test + void should_stream_chat_with_single_consumer_via_request() throws Exception { + + when(mockAuthenticator.tokenAsync()).thenReturn(completedFuture("token")); + wireMock.stubFor(post("/ml/v1/deployments/my-deployment-id/text/chat_stream?version=%s".formatted(API_VERSION)) + .withHeader("Authorization", equalTo("Bearer token")) + .willReturn(aResponse() + .withStatus(200) + .withBody( + """ + id: 1 + event: message + data: {"id":"chatcmpl-1","object":"chat.completion.chunk","model_id":"meta-llama/llama-4-maverick-17b-128e-instruct-fp8","model":"meta-llama/llama-4-maverick-17b-128e-instruct-fp8","choices":[{"index":0,"finish_reason":null,"delta":{"role":"assistant","content":"Cia"}}],"created":1749736055,"model_version":"4.0.0","created_at":"2025-06-12T13:47:35.541Z"} + + id: 2 + event: message + data: {"id":"chatcmpl-1","object":"chat.completion.chunk","model_id":"meta-llama/llama-4-maverick-17b-128e-instruct-fp8","model":"meta-llama/llama-4-maverick-17b-128e-instruct-fp8","choices":[{"index":0,"finish_reason":"stop","delta":{"content":"o"}}],"created":1749736055,"model_version":"4.0.0","created_at":"2025-06-12T13:47:35.552Z"} + """))); + + var deploymentService = DeploymentService.builder() + .baseUrl(URI.create("http://localhost:%s".formatted(wireMock.getPort()))) + .authenticator(mockAuthenticator) + .build(); + + var request = DeploymentChatRequest.builder() + .deploymentId("my-deployment-id") + .messages(UserMessage.text("Translate \"Hello\" in Italian")) + .build(); + + var response = new StringBuilder(); + deploymentService.chatStreaming(request, response::append).join(); + assertEquals("Ciao", response.toString()); + } + + @Test + void should_stream_thinking_and_response_via_consumers_from_request() throws Exception { + + String BODY = new String(ClassLoader.getSystemResourceAsStream("gpt_oss_thinking_streaming_response.txt").readAllBytes()); + + when(mockAuthenticator.tokenAsync()).thenReturn(completedFuture("token")); + wireMock.stubFor(post("/ml/v1/deployments/my-deployment-id/text/chat_stream?version=%s".formatted(API_VERSION)) + .withHeader("Authorization", equalTo("Bearer token")) + .willReturn(aResponse() + .withStatus(200) + .withChunkedDribbleDelay(29, 100) + .withBody(BODY))); + + var deploymentService = DeploymentService.builder() + .baseUrl(URI.create("http://localhost:%s".formatted(wireMock.getPort()))) + .authenticator(mockAuthenticator) + .build(); + + var request = DeploymentChatRequest.builder() + .deploymentId("my-deployment-id") + .messages(UserMessage.text("Translate \"Hello\" in Italian")) + .build(); + + var thinking = new StringBuilder(); + var response = new StringBuilder(); + deploymentService.chatStreaming(request, response::append, thinking::append).join(); + + assertEquals("User wants translation.", thinking.toString()); + assertTrue(response.toString().contains("\"ciao\"")); + } } diff --git a/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/chat/ChatServiceThinkingTest.java b/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/chat/ChatServiceThinkingTest.java index 10c118f9..d54ed2a3 100644 --- a/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/chat/ChatServiceThinkingTest.java +++ b/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/chat/ChatServiceThinkingTest.java @@ -1057,5 +1057,87 @@ public void onPartialThinking(String partialThinking, PartialChatResponse partia assertEquals(chatResponse.toAssistantMessage().thinking(), thinkingResponse.toString()); assertEquals(chatResponse.toAssistantMessage().content(), response.toString()); } + + @Test + void should_stream_thinking_and_response_via_string_consumers() throws Exception { + + var httpPort = wireMock.getPort(); + String BODY = new String(ClassLoader.getSystemResourceAsStream("granite_thinking_streaming_response.txt").readAllBytes()); + + wireMock.stubFor(post("/ml/v1/text/chat_stream?version=%s".formatted(API_VERSION)) + .withHeader("Authorization", equalTo("Bearer my-super-token")) + .willReturn(aResponse() + .withStatus(200) + .withChunkedDribbleDelay(159, 200) + .withBody(BODY))); + + when(mockAuthenticator.tokenAsync()).thenReturn(completedFuture("my-super-token")); + + var chatService = ChatService.builder() + .authenticator(mockAuthenticator) + .modelId("ibm/granite-3-3-8b-instruct") + .projectId("project-id") + .baseUrl(URI.create("http://localhost:%s".formatted(httpPort))) + .build(); + + var thinking = new StringBuilder(); + var response = new StringBuilder(); + + var chatRequest = ChatRequest.builder() + .messages(UserMessage.text("Translate \"Hello\" in Italian")) + .thinking(ExtractionTags.of(new Think("", ""), new Response("", ""))) + .build(); + + chatService.chatStreaming(chatRequest, response::append, thinking::append).join(); + + var EXPECTED_THINKING = + "The translation of \"Hello\" in Italian is straightforward. \"Hello\" in English directly translates to \"Ciao\" in Italian, which is a common informal greeting. For a more formal context, \"Buongiorno\" can be used, meaning \"Good day.\" However, since the request is for a direct translation of \"Hello,\" \"Ciao\" is the most appropriate response."; + var EXPECTED_RESPONSE = + "This is the informal equivalent, widely used in everyday conversation. For a formal greeting, one would say \"Buongiorno,\" but given the direct translation request, \"Ciao\" is the most fitting response."; + + assertTrue(thinking.toString().contains(EXPECTED_THINKING)); + assertTrue(response.toString().contains(EXPECTED_RESPONSE)); + } + + @Test + void should_stream_thinking_and_response_via_list_consumers() throws Exception { + + var httpPort = wireMock.getPort(); + String BODY = new String(ClassLoader.getSystemResourceAsStream("granite_thinking_streaming_response.txt").readAllBytes()); + + wireMock.stubFor(post("/ml/v1/text/chat_stream?version=%s".formatted(API_VERSION)) + .withHeader("Authorization", equalTo("Bearer my-super-token")) + .willReturn(aResponse() + .withStatus(200) + .withChunkedDribbleDelay(159, 200) + .withBody(BODY))); + + when(mockAuthenticator.tokenAsync()).thenReturn(completedFuture("my-super-token")); + + var chatService = ChatService.builder() + .authenticator(mockAuthenticator) + .modelId("ibm/granite-3-3-8b-instruct") + .projectId("project-id") + .baseUrl(URI.create("http://localhost:%s".formatted(httpPort))) + .build(); + + var thinking = new StringBuilder(); + var response = new StringBuilder(); + + var chatRequest = ChatRequest.builder() + .messages(UserMessage.text("Translate \"Hello\" in Italian")) + .thinking(ExtractionTags.of(new Think("", ""), new Response("", ""))) + .build(); + + chatService.chatStreaming(chatRequest, response::append, thinking::append).join(); + + var EXPECTED_THINKING = + "The translation of \"Hello\" in Italian is straightforward. \"Hello\" in English directly translates to \"Ciao\" in Italian, which is a common informal greeting. For a more formal context, \"Buongiorno\" can be used, meaning \"Good day.\" However, since the request is for a direct translation of \"Hello,\" \"Ciao\" is the most appropriate response."; + var EXPECTED_RESPONSE = + "This is the informal equivalent, widely used in everyday conversation. For a formal greeting, one would say \"Buongiorno,\" but given the direct translation request, \"Ciao\" is the most fitting response."; + + assertTrue(thinking.toString().contains(EXPECTED_THINKING)); + assertTrue(response.toString().contains(EXPECTED_RESPONSE)); + } } } diff --git a/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/client/impl/CustomProjectRestClient.java b/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/client/impl/CustomProjectRestClient.java index 3f5232c6..35199241 100644 --- a/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/client/impl/CustomProjectRestClient.java +++ b/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/client/impl/CustomProjectRestClient.java @@ -7,7 +7,6 @@ import java.util.Optional; import com.ibm.watsonx.ai.project.Project; import com.ibm.watsonx.ai.project.ProjectRestClient; -import com.ibm.watsonx.ai.project.ProjectRestClient.ProjectRestClientBuilderFactory; public class CustomProjectRestClient extends ProjectRestClient { @@ -20,8 +19,7 @@ public Optional findProject(String projectId) { throw new UnsupportedOperationException("Unimplemented method 'findProject'"); } - public static final class CustomProjectRestClientBuilderFactory - implements ProjectRestClientBuilderFactory { + public static final class CustomProjectRestClientBuilderFactory implements ProjectRestClientBuilderFactory { @Override public ProjectRestClient.Builder get() { return new CustomProjectRestClient.Builder(); diff --git a/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/gateway/chat/ModelGatewayChatServiceTest.java b/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/gateway/chat/ModelGatewayChatServiceTest.java index e9268fe4..f7fb88c4 100644 --- a/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/gateway/chat/ModelGatewayChatServiceTest.java +++ b/modules/watsonx-ai/src/test/java/com/ibm/watsonx/ai/gateway/chat/ModelGatewayChatServiceTest.java @@ -1048,4 +1048,72 @@ public void onError(Throwable error) { assertEquals("Ciao mondo!", completeResponse.toAssistantMessage().content()); assertEquals("Ciao mondo!", returnedResponse.toAssistantMessage().content()); } + + @Test + void should_stream_with_single_consumer_via_request() throws Exception { + + wireMock.stubFor(post("/ml/gateway/v1/chat/completions?version=%s".formatted(API_VERSION)) + .withHeader("Accept", equalTo("text/event-stream")) + .willReturn(aResponse() + .withStatus(200) + .withChunkedDribbleDelay(2, 100) + .withBody( + """ + data: {"id":"chatcmpl-r1","object":"chat.completion.chunk","choices":[{"index":0,"delta":{"role":"assistant","content":"Cia"},"finish_reason":"","logprobs":null}],"created":1,"model":"gpt-4o","usage":null,"cached":false} + + data: {"id":"chatcmpl-r1","object":"chat.completion.chunk","choices":[{"index":0,"delta":{"content":"o"},"finish_reason":"stop","logprobs":null}],"created":1,"model":"gpt-4o","usage":null,"cached":false} + + data: [DONE] + """))); + + when(mockAuthenticator.tokenAsync()).thenReturn(completedFuture("my-super-token")); + + var service = ModelGatewayChatService.builder() + .authenticator(mockAuthenticator) + .modelId("gpt-4o") + .baseUrl(URI.create("http://localhost:%s".formatted(wireMock.getPort()))) + .version(API_VERSION) + .build(); + + var request = ModelGatewayChatRequest.builder() + .messages(UserMessage.text("Translate \"Hello\" in Italian")) + .build(); + + var response = new StringBuilder(); + service.chatStreaming(request, response::append).join(); + assertEquals("Ciao", response.toString()); + } + + @Test + void should_stream_thinking_and_response_via_consumers_from_request() throws Exception { + + String BODY = new String(ClassLoader.getSystemResourceAsStream("gpt_oss_thinking_streaming_response.txt").readAllBytes()); + + wireMock.stubFor(post("/ml/gateway/v1/chat/completions?version=%s".formatted(API_VERSION)) + .withHeader("Accept", equalTo("text/event-stream")) + .willReturn(aResponse() + .withStatus(200) + .withChunkedDribbleDelay(29, 100) + .withBody(BODY))); + + when(mockAuthenticator.tokenAsync()).thenReturn(completedFuture("my-super-token")); + + var service = ModelGatewayChatService.builder() + .authenticator(mockAuthenticator) + .modelId("openai/gpt-oss-120b-curated") + .baseUrl(URI.create("http://localhost:%s".formatted(wireMock.getPort()))) + .version(API_VERSION) + .build(); + + var request = ModelGatewayChatRequest.builder() + .messages(UserMessage.text("Translate \"Hello\" in Italian")) + .build(); + + var thinking = new StringBuilder(); + var response = new StringBuilder(); + service.chatStreaming(request, response::append, thinking::append).join(); + + assertEquals("User wants translation.", thinking.toString()); + assertTrue(response.toString().contains("\"ciao\"")); + } } diff --git a/samples/pom.xml b/samples/pom.xml index feb9701b..a8a4e996 100644 --- a/samples/pom.xml +++ b/samples/pom.xml @@ -53,6 +53,12 @@ 0.40.0 + + com.ibm.watsonx + watsonx-ai-jackson2 + 0.40.0 + + io.smallrye.config smallrye-config-core