From d6ed1037e3de9140e30f87e86142efbeccef0de2 Mon Sep 17 00:00:00 2001 From: Dang Zitou Date: Thu, 13 Aug 2026 03:39:07 +0800 Subject: [PATCH 1/2] Fix publish event content type handling Signed-off-by: Dang Zitou --- .../java/io/dapr/client/DaprClientImpl.java | 34 ++++++-- .../io/dapr/client/DaprClientGrpcTest.java | 83 +++++++++++++++---- 2 files changed, 93 insertions(+), 24 deletions(-) diff --git a/sdk/src/main/java/io/dapr/client/DaprClientImpl.java b/sdk/src/main/java/io/dapr/client/DaprClientImpl.java index de91813c50..39962462d5 100644 --- a/sdk/src/main/java/io/dapr/client/DaprClientImpl.java +++ b/sdk/src/main/java/io/dapr/client/DaprClientImpl.java @@ -352,18 +352,20 @@ public Mono publishEvent(PublishEventRequest request) { String pubsubName = request.getPubsubName(); String topic = request.getTopic(); Object data = request.getData(); + String contentType = getPublishEventContentType(data, request.getContentType()); + boolean useContentTypeConverter = objectSerializer instanceof DefaultObjectSerializer + && (DefaultContentTypeConverter.isBinaryContentType(contentType) + || DefaultContentTypeConverter.isStringContentType(contentType) + || DefaultContentTypeConverter.isJsonContentType(contentType) + || DefaultContentTypeConverter.isCloudEventContentType(contentType)); + byte[] serializedEvent = useContentTypeConverter + ? DefaultContentTypeConverter.convertEventToBytesForGrpc(data, contentType) + : objectSerializer.serialize(data); DaprPubsubProtos.PublishEventRequest.Builder envelopeBuilder = DaprPubsubProtos.PublishEventRequest.newBuilder() .setTopic(topic) .setPubsubName(pubsubName) - .setData(ByteString.copyFrom(objectSerializer.serialize(data))); - - // Content-type can be overwritten on a per-request basis. - // It allows CloudEvents to be handled differently, for example. - String contentType = request.getContentType(); - if (contentType == null || contentType.isEmpty()) { - contentType = objectSerializer.getContentType(); - } - envelopeBuilder.setDataContentType(contentType); + .setData(ByteString.copyFrom(serializedEvent)) + .setDataContentType(contentType); Map metadata = request.getMetadata(); if (metadata != null) { @@ -381,6 +383,20 @@ public Mono publishEvent(PublishEventRequest request) { } } + private String getPublishEventContentType(Object data, String contentType) { + if (!Strings.isNullOrEmpty(contentType) || !(objectSerializer instanceof DefaultObjectSerializer)) { + return Strings.isNullOrEmpty(contentType) ? objectSerializer.getContentType() : contentType; + } + + if (data instanceof byte[]) { + return "application/octet-stream"; + } + if (data instanceof String || data instanceof Boolean || data instanceof Number) { + return "text/plain"; + } + return objectSerializer.getContentType(); + } + /** * {@inheritDoc} */ diff --git a/sdk/src/test/java/io/dapr/client/DaprClientGrpcTest.java b/sdk/src/test/java/io/dapr/client/DaprClientGrpcTest.java index 6af09aa190..55ef71897f 100644 --- a/sdk/src/test/java/io/dapr/client/DaprClientGrpcTest.java +++ b/sdk/src/test/java/io/dapr/client/DaprClientGrpcTest.java @@ -231,27 +231,80 @@ public void publishEventObjectTest() { @Test public void publishEventContentTypeOverrideTest() { + mockPublishEventSuccess(); + + Mono result = client.publishEvent( + new PublishEventRequest("pubsubname", "topic", "hello") + .setContentType("text/plain")); + result.block(); + + DaprPubsubProtos.PublishEventRequest request = capturePublishEventRequest(); + assertEquals("text/plain", request.getDataContentType()); + assertEquals("hello", request.getData().toString(StandardCharsets.UTF_8)); + } + + @Test + public void publishEventStringInfersTextPlainContentTypeTest() { + mockPublishEventSuccess(); + + client.publishEvent("pubsubname", "topic", "hello").block(); + + DaprPubsubProtos.PublishEventRequest request = capturePublishEventRequest(); + assertEquals("text/plain", request.getDataContentType()); + assertEquals("hello", request.getData().toString(StandardCharsets.UTF_8)); + } + + @Test + public void publishEventNumberInfersTextPlainContentTypeTest() { + mockPublishEventSuccess(); + + client.publishEvent("pubsubname", "topic", 42).block(); + + DaprPubsubProtos.PublishEventRequest request = capturePublishEventRequest(); + assertEquals("text/plain", request.getDataContentType()); + assertEquals("42", request.getData().toString(StandardCharsets.UTF_8)); + } + + @Test + public void publishEventByteArrayInfersOctetStreamContentTypeTest() { + byte[] event = new byte[] {1, 2, 3}; + mockPublishEventSuccess(); + + client.publishEvent("pubsubname", "topic", event).block(); + + DaprPubsubProtos.PublishEventRequest request = capturePublishEventRequest(); + assertEquals("application/octet-stream", request.getDataContentType()); + assertArrayEquals(event, request.getData().toByteArray()); + } + + @Test + public void publishEventByteArrayPreservesCustomContentTypeTest() { + byte[] event = new byte[] {1, 2, 3}; + mockPublishEventSuccess(); + + client.publishEvent( + new PublishEventRequest("pubsubname", "topic", event) + .setContentType("image/png")).block(); + + DaprPubsubProtos.PublishEventRequest request = capturePublishEventRequest(); + assertEquals("image/png", request.getDataContentType()); + assertArrayEquals(event, request.getData().toByteArray()); + } + + private void mockPublishEventSuccess() { doAnswer((Answer) invocation -> { StreamObserver observer = (StreamObserver) invocation.getArguments()[1]; observer.onNext(Empty.getDefaultInstance()); observer.onCompleted(); return null; - }).when(daprStub).publishEvent(ArgumentMatchers.argThat(publishEventRequest -> { - if (!"text/plain".equals(publishEventRequest.getDataContentType())) { - return false; - } - - if (!"\"hello\"".equals(new String(publishEventRequest.getData().toByteArray()))) { - return false; - } - return true; - }), any()); - + }).when(daprStub).publishEvent(any(), any()); + } - Mono result = client.publishEvent( - new PublishEventRequest("pubsubname", "topic", "hello") - .setContentType("text/plain")); - result.block(); + private DaprPubsubProtos.PublishEventRequest capturePublishEventRequest() { + ArgumentCaptor captor = + ArgumentCaptor.forClass(DaprPubsubProtos.PublishEventRequest.class); + verify(daprStub).publishEvent(captor.capture(), any()); + return captor.getValue(); } @Test From aab88d6299be3d8fa66a6bbd13b0adbf437963b6 Mon Sep 17 00:00:00 2001 From: Dang Zitou Date: Mon, 17 Aug 2026 01:29:19 +0800 Subject: [PATCH 2/2] chore(ci): rerun checks Signed-off-by: Dang Zitou