This is an automated email from the ASF dual-hosted git repository.

fjtiradosarti pushed a commit to branch main
in repository 
https://gitbox.apache.org/repos/asf/incubator-kie-kogito-runtimes.git


The following commit(s) were added to refs/heads/main by this push:
     new b8b74df59c [Fix #3763] Properly reading compressdata string header 
(#3764)
b8b74df59c is described below

commit b8b74df59c2c54c34fb049ecb7bc341853842991
Author: Francisco Javier Tirado Sarti 
<[email protected]>
AuthorDate: Wed Nov 6 14:41:26 2024 +0100

    [Fix #3763] Properly reading compressdata string header (#3764)
    
    * [Fix #3763] Properly reading compressdata string header
    
    * [Fix #3763] Adding unit test
---
 .../process/MultipleProcessInstanceDataEvent.java  |  12 +-
 ...ultipleProcessDataInstanceConverterFactory.java |  22 +++-
 ...MultipleProcessInstanceDataEventSerializer.java |   4 +-
 .../kogito/event/process/ProcessEventsTest.java    | 126 ++++++++++++++++++---
 4 files changed, 139 insertions(+), 25 deletions(-)

diff --git 
a/api/kogito-events-core/src/main/java/org/kie/kogito/event/process/MultipleProcessInstanceDataEvent.java
 
b/api/kogito-events-core/src/main/java/org/kie/kogito/event/process/MultipleProcessInstanceDataEvent.java
index f29a920c13..a2783038b2 100644
--- 
a/api/kogito-events-core/src/main/java/org/kie/kogito/event/process/MultipleProcessInstanceDataEvent.java
+++ 
b/api/kogito-events-core/src/main/java/org/kie/kogito/event/process/MultipleProcessInstanceDataEvent.java
@@ -35,8 +35,16 @@ public class MultipleProcessInstanceDataEvent extends 
ProcessInstanceDataEvent<C
     }
 
     public boolean isCompressed() {
-        Object extension = 
getExtension(MultipleProcessInstanceDataEvent.COMPRESS_DATA);
-        return extension instanceof Boolean ? ((Boolean) 
extension).booleanValue() : false;
+        return 
isCompressed(getExtension(MultipleProcessInstanceDataEvent.COMPRESS_DATA));
+    }
+
+    public static boolean isCompressed(Object extension) {
+        if (extension instanceof Boolean) {
+            return ((Boolean) extension).booleanValue();
+        } else if (extension instanceof String) {
+            return Boolean.parseBoolean((String) extension);
+        }
+        return false;
     }
 
     public void setCompressed(boolean compressed) {
diff --git 
a/api/kogito-events-core/src/main/java/org/kie/kogito/event/serializer/MultipleProcessDataInstanceConverterFactory.java
 
b/api/kogito-events-core/src/main/java/org/kie/kogito/event/serializer/MultipleProcessDataInstanceConverterFactory.java
index 4f8a6205c2..0c7dd53901 100644
--- 
a/api/kogito-events-core/src/main/java/org/kie/kogito/event/serializer/MultipleProcessDataInstanceConverterFactory.java
+++ 
b/api/kogito-events-core/src/main/java/org/kie/kogito/event/serializer/MultipleProcessDataInstanceConverterFactory.java
@@ -33,12 +33,20 @@ import com.fasterxml.jackson.databind.ObjectMapper;
 
 import io.cloudevents.CloudEvent;
 import io.cloudevents.CloudEventData;
+import io.cloudevents.core.data.PojoCloudEventData.ToBytes;
 
 public class MultipleProcessDataInstanceConverterFactory {
-
     private MultipleProcessDataInstanceConverterFactory() {
     }
 
+    public static ToBytes<Collection<ProcessInstanceDataEvent<? extends 
KogitoMarshallEventSupport>>> toCloudEvent(MultipleProcessInstanceDataEvent 
event, ObjectMapper objectMapper) {
+        if 
(MultipleProcessInstanceDataEvent.BINARY_CONTENT_TYPE.equals(event.getDataContentType()))
 {
+            return event.isCompressed() ? compressedToBytes : binaryToBytes;
+        } else {
+            return objectMapper::writeValueAsBytes;
+        }
+    }
+
     public static Converter<CloudEventData, 
Collection<ProcessInstanceDataEvent<? extends KogitoMarshallEventSupport>>> 
fromCloudEvent(CloudEvent cloudEvent, ObjectMapper objectMapper) {
         if 
(MultipleProcessInstanceDataEvent.BINARY_CONTENT_TYPE.equals(cloudEvent.getDataContentType()))
 {
             return isCompressed(cloudEvent) ? compressedConverter : 
binaryConverter;
@@ -49,10 +57,13 @@ public class MultipleProcessDataInstanceConverterFactory {
     }
 
     private static boolean isCompressed(CloudEvent event) {
-        Object value = 
event.getExtension(MultipleProcessInstanceDataEvent.COMPRESS_DATA);
-        return value instanceof Boolean ? ((Boolean) value).booleanValue() : 
false;
+        return 
MultipleProcessInstanceDataEvent.isCompressed(event.getExtension(MultipleProcessInstanceDataEvent.COMPRESS_DATA));
     }
 
+    private static ToBytes<Collection<ProcessInstanceDataEvent<? extends 
KogitoMarshallEventSupport>>> compressedToBytes = data -> serialize(data, true);
+
+    private static ToBytes<Collection<ProcessInstanceDataEvent<? extends 
KogitoMarshallEventSupport>>> binaryToBytes = data -> serialize(data, false);
+
     private static Converter<CloudEventData, 
Collection<ProcessInstanceDataEvent<? extends KogitoMarshallEventSupport>>> 
binaryConverter =
             data -> deserialize(data, false);
 
@@ -62,4 +73,9 @@ public class MultipleProcessDataInstanceConverterFactory {
     private static Collection<ProcessInstanceDataEvent<? extends 
KogitoMarshallEventSupport>> deserialize(CloudEventData data, boolean compress) 
throws IOException {
         return 
MultipleProcessInstanceDataEventDeserializer.readFromBytes(Base64.getDecoder().decode(data.toBytes()),
 compress);
     }
+
+    private static byte[] serialize(Collection<ProcessInstanceDataEvent<? 
extends KogitoMarshallEventSupport>> data,
+            boolean compress) throws IOException {
+        return 
Base64.getEncoder().encode(MultipleProcessInstanceDataEventSerializer.dataAsBytes(data,
 compress));
+    }
 }
diff --git 
a/api/kogito-events-core/src/main/java/org/kie/kogito/event/serializer/MultipleProcessInstanceDataEventSerializer.java
 
b/api/kogito-events-core/src/main/java/org/kie/kogito/event/serializer/MultipleProcessInstanceDataEventSerializer.java
index 42825e9679..42b219b46f 100644
--- 
a/api/kogito-events-core/src/main/java/org/kie/kogito/event/serializer/MultipleProcessInstanceDataEventSerializer.java
+++ 
b/api/kogito-events-core/src/main/java/org/kie/kogito/event/serializer/MultipleProcessInstanceDataEventSerializer.java
@@ -62,14 +62,14 @@ public class MultipleProcessInstanceDataEventSerializer 
extends JsonSerializer<M
             if (compress) {
                 
gen.writeBooleanField(MultipleProcessInstanceDataEvent.COMPRESS_DATA, true);
             }
-            gen.writeBinaryField("data", dataAsBytes(gen, value.getData(), 
compress));
+            gen.writeBinaryField("data", dataAsBytes(value.getData(), 
compress));
             gen.writeEndObject();
         } else {
             defaultSerializer.serialize(value, gen, serializers);
         }
     }
 
-    private byte[] dataAsBytes(JsonGenerator gen, 
Collection<ProcessInstanceDataEvent<? extends KogitoMarshallEventSupport>> 
data, boolean compress) throws IOException {
+    static byte[] dataAsBytes(Collection<ProcessInstanceDataEvent<? extends 
KogitoMarshallEventSupport>> data, boolean compress) throws IOException {
         ByteArrayOutputStream bytesOut = new ByteArrayOutputStream();
         try (DataOutputStream out = new DataOutputStream(compress ? new 
GZIPOutputStream(bytesOut) : bytesOut)) {
             logger.trace("Writing size {}", data.size());
diff --git 
a/api/kogito-events-core/src/test/java/org/kie/kogito/event/process/ProcessEventsTest.java
 
b/api/kogito-events-core/src/test/java/org/kie/kogito/event/process/ProcessEventsTest.java
index d278e2c427..50d2f6e867 100644
--- 
a/api/kogito-events-core/src/test/java/org/kie/kogito/event/process/ProcessEventsTest.java
+++ 
b/api/kogito-events-core/src/test/java/org/kie/kogito/event/process/ProcessEventsTest.java
@@ -22,8 +22,11 @@ import java.io.IOException;
 import java.net.URI;
 import java.time.OffsetDateTime;
 import java.util.Arrays;
+import java.util.HashMap;
 import java.util.Iterator;
+import java.util.Map;
 import java.util.Set;
+import java.util.function.BiConsumer;
 import java.util.stream.Collectors;
 
 import org.junit.jupiter.api.Test;
@@ -40,8 +43,17 @@ import com.fasterxml.jackson.databind.JsonNode;
 import com.fasterxml.jackson.databind.ObjectMapper;
 
 import io.cloudevents.CloudEvent;
+import io.cloudevents.CloudEventData;
 import io.cloudevents.SpecVersion;
+import io.cloudevents.core.data.BytesCloudEventData;
+import io.cloudevents.core.data.PojoCloudEventData;
+import io.cloudevents.core.format.EventFormat;
+import io.cloudevents.core.message.MessageWriter;
+import io.cloudevents.core.message.impl.BaseGenericBinaryMessageReaderImpl;
 import io.cloudevents.jackson.JsonFormat;
+import io.cloudevents.rw.CloudEventContextWriter;
+import io.cloudevents.rw.CloudEventRWException;
+import io.cloudevents.rw.CloudEventWriter;
 
 import static org.assertj.core.api.Assertions.assertThat;
 import static 
org.kie.kogito.event.process.KogitoEventBodySerializationHelper.toDate;
@@ -130,14 +142,103 @@ class ProcessEventsTest {
     @Test
     void multipleInstanceDataEvent() throws IOException {
         JsonNode expectedVarValue = 
OBJECT_MAPPER.createObjectNode().put("name", "John Doe");
-        int standard = processMultipleInstanceDataEvent(expectedVarValue, 
false, false);
-        int binary = processMultipleInstanceDataEvent(expectedVarValue, true, 
false);
-        int binaryCompressed = 
processMultipleInstanceDataEvent(expectedVarValue, true, true);
-        assertThat(standard).isGreaterThan(binary);
-        assertThat(binary).isGreaterThan(binaryCompressed);
+        processMultipleInstanceDataEvent(expectedVarValue, false, false, 
this::serializeAsStructured);
+        processMultipleInstanceDataEvent(expectedVarValue, true, false, 
this::serializeAsStructured);
+        processMultipleInstanceDataEvent(expectedVarValue, true, true, 
this::serializeAsStructured);
+        processMultipleInstanceDataEvent(expectedVarValue, false, false, 
this::serializeAsBinary);
+        processMultipleInstanceDataEvent(expectedVarValue, true, false, 
this::serializeAsBinary);
+        processMultipleInstanceDataEvent(expectedVarValue, true, true, 
this::serializeAsBinary);
     }
 
-    private int processMultipleInstanceDataEvent(JsonNode expectedVarValue, 
boolean binary, boolean compress) throws IOException {
+    private MultipleProcessInstanceDataEvent 
serializeAsStructured(MultipleProcessInstanceDataEvent event) throws 
IOException {
+        return OBJECT_MAPPER.readValue(OBJECT_MAPPER.writeValueAsBytes(event), 
MultipleProcessInstanceDataEvent.class);
+    }
+
+    private record CloudEventHolder(Map<String, String> headers, byte[] data) {
+    }
+
+    private static class TestMessageWriter implements 
CloudEventWriter<CloudEventHolder>, 
MessageWriter<CloudEventWriter<CloudEventHolder>, CloudEventHolder> {
+
+        private Map<String, String> headers = new HashMap<>();
+        private byte[] value;
+
+        @Override
+        public TestMessageWriter create(SpecVersion version) throws 
CloudEventRWException {
+            headers.put("specversion", version.toString());
+            return this;
+        }
+
+        @Override
+        public CloudEventHolder setEvent(EventFormat format, byte[] value) 
throws CloudEventRWException {
+            this.value = value;
+            return this.end();
+        }
+
+        @Override
+        public CloudEventContextWriter withContextAttribute(String name, 
String value) throws CloudEventRWException {
+            headers.put(name, value);
+            return this;
+        }
+
+        @Override
+        public CloudEventHolder end(CloudEventData data) throws 
CloudEventRWException {
+            this.value = data.toBytes();
+            return this.end();
+        }
+
+        @Override
+        public CloudEventHolder end() throws CloudEventRWException {
+            return new CloudEventHolder(headers, value);
+        }
+    }
+
+    private class TestMessageReader extends 
BaseGenericBinaryMessageReaderImpl<String, String> {
+        private Map<String, String> headers;
+
+        protected TestMessageReader(SpecVersion version, CloudEventHolder 
body) {
+            super(version, body.data() == null ? null : 
BytesCloudEventData.wrap(body.data()));
+            this.headers = body.headers();
+        }
+
+        @Override
+        protected boolean isContentTypeHeader(String key) {
+            return false;
+        }
+
+        @Override
+        protected boolean isCloudEventsHeader(String key) {
+            return true;
+        }
+
+        @Override
+        protected String toCloudEventsKey(String key) {
+            return key;
+        }
+
+        @Override
+        protected void forEachHeader(BiConsumer<String, String> fn) {
+            headers.forEach(fn);
+        }
+
+        @Override
+        protected String toCloudEventsValue(String value) {
+            return value;
+        }
+
+    }
+
+    private MultipleProcessInstanceDataEvent 
serializeAsBinary(MultipleProcessInstanceDataEvent event) throws IOException {
+        CloudEvent toSerialize = event.asCloudEvent(value -> 
PojoCloudEventData.wrap(value, 
MultipleProcessDataInstanceConverterFactory.toCloudEvent(event, 
OBJECT_MAPPER)));
+        CloudEventHolder holder = new 
TestMessageWriter().writeBinary(toSerialize);
+        CloudEvent deserialized = new TestMessageReader(SpecVersion.V1, 
holder).toEvent();
+        return DataEventFactory.from(new MultipleProcessInstanceDataEvent(), 
deserialized, 
MultipleProcessDataInstanceConverterFactory.fromCloudEvent(deserialized, 
OBJECT_MAPPER));
+    }
+
+    private static interface CheckedUnaryOperator<T> {
+        T apply(T obj) throws IOException;
+    }
+
+    private void processMultipleInstanceDataEvent(JsonNode expectedVarValue, 
boolean binary, boolean compress, 
CheckedUnaryOperator<MultipleProcessInstanceDataEvent> operator) throws 
IOException {
         ProcessInstanceStateDataEvent stateEvent = new 
ProcessInstanceStateDataEvent();
         setBaseEventValues(stateEvent, 
ProcessInstanceStateDataEvent.STATE_TYPE);
         
stateEvent.setData(ProcessInstanceStateEventBody.create().eventDate(toDate(TIME)).eventType(EVENT_TYPE).eventUser(SUBJECT)
@@ -185,20 +286,9 @@ class ProcessEventsTest {
             event.setCompressed(compress);
         }
 
-        byte[] json = OBJECT_MAPPER.writeValueAsBytes(event);
-        logger.info("Serialized chunk size is {}", json.length);
-
-        // cloud event structured mode check
-        MultipleProcessInstanceDataEvent deserializedEvent = 
OBJECT_MAPPER.readValue(json, MultipleProcessInstanceDataEvent.class);
-        
assertThat(deserializedEvent.getData()).hasSize(event.getData().size());
-        assertMultipleIntance(deserializedEvent, expectedVarValue);
-
-        // cloud event binary mode check
-        CloudEvent cloudEvent = OBJECT_MAPPER.readValue(json, 
CloudEvent.class);
-        deserializedEvent = DataEventFactory.from(new 
MultipleProcessInstanceDataEvent(), cloudEvent, 
MultipleProcessDataInstanceConverterFactory.fromCloudEvent(cloudEvent, 
OBJECT_MAPPER));
+        MultipleProcessInstanceDataEvent deserializedEvent = 
operator.apply(event);
         
assertThat(deserializedEvent.getData()).hasSize(event.getData().size());
         assertMultipleIntance(deserializedEvent, expectedVarValue);
-        return json.length;
     }
 
     private void assertMultipleIntance(MultipleProcessInstanceDataEvent 
deserializedEvent, JsonNode expectedVarValue) {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to