GH-924 Fix regression with structured CE cnversion into Message
Resolves #924
This commit is contained in:
@@ -231,15 +231,16 @@ public final class CloudEventMessageUtils {
|
|||||||
inputMessage = canonicalizeHeadersWithPossibleCopy(inputMessage);
|
inputMessage = canonicalizeHeadersWithPossibleCopy(inputMessage);
|
||||||
Map<String, Object> headers = new HashMap<>(inputMessage.getHeaders());
|
Map<String, Object> headers = new HashMap<>(inputMessage.getHeaders());
|
||||||
|
|
||||||
if (isCloudEvent(inputMessage) && headers.containsKey("content-type")) {
|
boolean isCloudEvent = isCloudEvent(inputMessage);
|
||||||
|
if (isCloudEvent && headers.containsKey("content-type")) {
|
||||||
inputMessage = MessageBuilder.fromMessage(inputMessage).setHeader(MessageHeaders.CONTENT_TYPE, headers.get("content-type")).build();
|
inputMessage = MessageBuilder.fromMessage(inputMessage).setHeader(MessageHeaders.CONTENT_TYPE, headers.get("content-type")).build();
|
||||||
}
|
}
|
||||||
|
MimeType contentType = contentTypeResolver.resolve(inputMessage.getHeaders());
|
||||||
String inputContentType = (String) inputMessage.getHeaders().get(DATACONTENTTYPE);
|
String inputContentType = (String) inputMessage.getHeaders().get(DATACONTENTTYPE);
|
||||||
// first check the obvious and see if content-type is `cloudevents`
|
// first check the obvious and see if content-type is `cloudevents`
|
||||||
if (!isCloudEvent(inputMessage) && headers.containsKey(MessageHeaders.CONTENT_TYPE)) {
|
if (!isCloudEvent && contentType != null) {
|
||||||
// structured-mode
|
// structured-mode
|
||||||
MimeType contentType = contentTypeResolver.resolve(inputMessage.getHeaders());
|
|
||||||
if (contentType.getType().equals(APPLICATION_CLOUDEVENTS.getType()) && contentType
|
if (contentType.getType().equals(APPLICATION_CLOUDEVENTS.getType()) && contentType
|
||||||
.getSubtype().startsWith(APPLICATION_CLOUDEVENTS.getSubtype())) {
|
.getSubtype().startsWith(APPLICATION_CLOUDEVENTS.getSubtype())) {
|
||||||
|
|
||||||
@@ -254,7 +255,6 @@ public final class CloudEventMessageUtils {
|
|||||||
.setHeader(DATACONTENTTYPE, dataContentType).build();
|
.setHeader(DATACONTENTTYPE, dataContentType).build();
|
||||||
Map<String, Object> structuredCloudEvent = (Map<String, Object>) messageConverter
|
Map<String, Object> structuredCloudEvent = (Map<String, Object>) messageConverter
|
||||||
.fromMessage(cloudEventMessage, Map.class);
|
.fromMessage(cloudEventMessage, Map.class);
|
||||||
|
|
||||||
canonicalizeHeaders(structuredCloudEvent, true);
|
canonicalizeHeaders(structuredCloudEvent, true);
|
||||||
return buildBinaryMessageFromStructuredMap(structuredCloudEvent,
|
return buildBinaryMessageFromStructuredMap(structuredCloudEvent,
|
||||||
inputMessage.getHeaders());
|
inputMessage.getHeaders());
|
||||||
|
|||||||
@@ -391,6 +391,7 @@ public class CloudEventFunctionTests {
|
|||||||
Function<Message<SpringReleaseEvent>, Message<SpringReleaseEvent>> springReleaseAsMessage() {
|
Function<Message<SpringReleaseEvent>, Message<SpringReleaseEvent>> springReleaseAsMessage() {
|
||||||
return message -> {
|
return message -> {
|
||||||
SpringReleaseEvent updated = springRelease().apply(message.getPayload());
|
SpringReleaseEvent updated = springRelease().apply(message.getPayload());
|
||||||
|
assertThat(message.getHeaders().get("ce-type")).isEqualTo("org.springframework");
|
||||||
return CloudEventMessageBuilder.withData(updated)
|
return CloudEventMessageBuilder.withData(updated)
|
||||||
.copyHeaders(message.getHeaders())
|
.copyHeaders(message.getHeaders())
|
||||||
.setSource("https://spring.release.event")
|
.setSource("https://spring.release.event")
|
||||||
|
|||||||
Reference in New Issue
Block a user