diff --git a/agentscope-extensions/agentscope-extensions-model/agentscope-extensions-model-dashscope/src/main/java/io/agentscope/extensions/model/dashscope/DashScopeHttpClient.java b/agentscope-extensions/agentscope-extensions-model/agentscope-extensions-model-dashscope/src/main/java/io/agentscope/extensions/model/dashscope/DashScopeHttpClient.java index 6adce7e2f2..613a0776bb 100644 --- a/agentscope-extensions/agentscope-extensions-model/agentscope-extensions-model-dashscope/src/main/java/io/agentscope/extensions/model/dashscope/DashScopeHttpClient.java +++ b/agentscope-extensions/agentscope-extensions-model/agentscope-extensions-model-dashscope/src/main/java/io/agentscope/extensions/model/dashscope/DashScopeHttpClient.java @@ -270,33 +270,36 @@ public Flux stream( .build(); return transport.stream(httpRequest) - .map( - data -> { + .handle( + (data, sink) -> { try { // Decrypt response if encryption is enabled if (finalEncryptionContext != null) { data = decryptResponse(data, finalEncryptionContext); } - return JsonUtils.getJsonCodec() - .fromJson(data, DashScopeResponse.class); + sink.next( + new ParsedStreamResponse( + data, + JsonUtils.getJsonCodec() + .fromJson( + data, + DashScopeResponse.class))); } catch (JsonException e) { log.warn( "Failed to parse SSE data: {}. Error: {}", data, e.getMessage()); - // Return null and filter out later - return null; } }) - .filter(response -> response != null) .handle( - (response, sink) -> { + (streamResponse, sink) -> { + DashScopeResponse response = streamResponse.response(); if (response.isError()) { sink.error( new DashScopeHttpException( "DashScope API error: " + response.getMessage(), response.getCode(), - null)); + streamResponse.responseBody())); } else { sink.next(response); } @@ -834,6 +837,8 @@ public DashScopeHttpClient build() { } } + private record ParsedStreamResponse(String responseBody, DashScopeResponse response) {} + /** * Exception thrown when DashScope HTTP operations fail. */ diff --git a/agentscope-extensions/agentscope-extensions-model/agentscope-extensions-model-dashscope/src/test/java/io/agentscope/extensions/model/dashscope/DashScopeHttpClientTest.java b/agentscope-extensions/agentscope-extensions-model/agentscope-extensions-model-dashscope/src/test/java/io/agentscope/extensions/model/dashscope/DashScopeHttpClientTest.java index 0d327bf690..30515f0677 100644 --- a/agentscope-extensions/agentscope-extensions-model/agentscope-extensions-model-dashscope/src/test/java/io/agentscope/extensions/model/dashscope/DashScopeHttpClientTest.java +++ b/agentscope-extensions/agentscope-extensions-model/agentscope-extensions-model-dashscope/src/test/java/io/agentscope/extensions/model/dashscope/DashScopeHttpClientTest.java @@ -647,10 +647,26 @@ void testStreamErrorHandling() { && dashScopeHttpException.getErrorCode().equals(errorCode) && dashScopeHttpException .getMessage() - .equals("DashScope API error: " + errorMessage)) + .equals("DashScope API error: " + errorMessage) + && dashScopeHttpException + .getResponseBody() + .contains("\"request_id\":\"request_id_123\"")) .verify(); } + @Test + void testStreamIgnoresMalformedSseData() { + mockServer.enqueue( + new MockResponse() + .setResponseCode(200) + .setBody("data: malformed-json\\n\\n") + .setHeader("Content-Type", "text/event-stream")); + + DashScopeRequest request = createTestRequest("qwen-plus", "test"); + + StepVerifier.create(client.stream(request, null, null, null)).verifyComplete(); + } + @Test void testHeaderOverride() throws Exception { mockServer.enqueue(