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]