From 9bed4da212aeeeecbaa35b576001b88d2be8fe6e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rafa=C5=82=20G=C5=82owacz?= Date: Tue, 21 Oct 2025 12:26:24 +0200 Subject: [PATCH 1/2] DXP-2378 StreamX rendering service does not start with RocksDB enabled --- .../ProcessRenderingRequestFunction.java | 2 +- .../engine/ProcessTriggersFunctions.java | 9 +++++++-- .../rendering/engine/RenderingRequests.java | 17 ++++++++++++++--- .../PreservedDataMessageConverter.java | 18 ++++++++++++------ ...servedRenderingContextMessageConverter.java | 18 ++++++++++++------ 5 files changed, 46 insertions(+), 18 deletions(-) diff --git a/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/ProcessRenderingRequestFunction.java b/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/ProcessRenderingRequestFunction.java index 78b2672d5..d4a317af9 100644 --- a/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/ProcessRenderingRequestFunction.java +++ b/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/ProcessRenderingRequestFunction.java @@ -79,7 +79,7 @@ public GenericPayload process(RenderingRequest request, Action request Action outputAction = getOutputAction(requestAction, Action.from(data.getMetadata()), Action.from(renderer.getMetadata())); - if (renderer.getPayload() != null || !PUBLISH.equals(outputAction)) { + if (renderer.getPayload() != null || UNPUBLISH.equals(outputAction)) { return getOutput(outputAction, request.getOutputKeyTemplate(), request.getOutputTypeTemplate(), request.getOutputFormat(), getValue(data), renderer.getPayload()); diff --git a/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/ProcessTriggersFunctions.java b/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/ProcessTriggersFunctions.java index b7fa8f805..53b9f2bde 100644 --- a/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/ProcessTriggersFunctions.java +++ b/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/ProcessTriggersFunctions.java @@ -14,6 +14,7 @@ import jakarta.enterprise.inject.Produces; import jakarta.inject.Inject; import java.util.List; +import java.util.Optional; import org.apache.pulsar.client.api.Schema; import org.eclipse.microprofile.reactive.messaging.Incoming; import org.eclipse.microprofile.reactive.messaging.Message; @@ -85,8 +86,12 @@ public Multi> processRenderer(Message incomi @Outgoing(Channels.Outgoing.RENDERING_REQUESTS) public Multi> processContext( Message incoming) { - RenderingContext renderingContext = incoming.getPayload().getRenderingContext(); - if (renderingContexts.hasRenderer(renderingContext)) { + boolean hasRenderer = Optional.ofNullable(incoming.getPayload()) + .map(PreservedRenderingContext::getRenderingContext) + .map(renderingContext -> renderingContexts.hasRenderer(renderingContext)) + .orElse(false); + if (hasRenderer) { + RenderingContext renderingContext = incoming.getPayload().getRenderingContext(); return renderingRequests.getFromDataStore(incoming, List.of(new KeyedValue<>(extractKey(incoming), renderingContext))); } else { diff --git a/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/RenderingRequests.java b/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/RenderingRequests.java index 4550d65bc..c11d8bd87 100644 --- a/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/RenderingRequests.java +++ b/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/RenderingRequests.java @@ -1,6 +1,7 @@ package dev.streamx.blueprints.rendering.engine; import static dev.streamx.blueprints.rendering.engine.RenderingContexts.isMatchingData; +import static dev.streamx.quasar.reactive.messaging.metadata.Action.PUBLISH; import static dev.streamx.quasar.reactive.messaging.utils.MetadataUtils.extractAction; import static dev.streamx.quasar.reactive.messaging.utils.MetadataUtils.extractEventTime; import static dev.streamx.quasar.reactive.messaging.utils.MetadataUtils.extractKey; @@ -18,6 +19,8 @@ import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; import java.util.List; +import java.util.Objects; +import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.function.Supplier; import java.util.stream.Stream; @@ -115,12 +118,20 @@ private Stream> calculateRenderingRequestForDataStore( private Stream> fetchStoredDataContextFromStore() { return dataStore.entriesWithMetadata() - .map(entry -> new KeyedValue<>(entry.key(), entry.value().getPayload().getData())); + .filter(entry -> + PUBLISH.equals(Action.from(entry.value().getMetadata()))) + .map(entry -> + new KeyedValue<>(entry.key(), Optional.ofNullable(entry.value().getPayload()) + .map(PreservedData::getData) + .orElse(null))) + .filter(entry -> entry.value() != null); } private boolean skipDataWithNoValue(KeyedValue entry) { - return entry.value() != null - || dataStore.get(entry.key()).getData() != null; + return Objects.nonNull(entry.value()) + || Optional.ofNullable(entry.key()) + .map(dataStore::get) + .map(PreservedData::getData).isPresent(); } private String getDataType(String key) { diff --git a/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/converter/PreservedDataMessageConverter.java b/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/converter/PreservedDataMessageConverter.java index e0943278f..93be18c54 100644 --- a/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/converter/PreservedDataMessageConverter.java +++ b/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/converter/PreservedDataMessageConverter.java @@ -11,6 +11,7 @@ import io.smallrye.reactive.messaging.MessageConverter; import jakarta.enterprise.context.ApplicationScoped; import java.lang.reflect.Type; +import java.util.Optional; import org.eclipse.microprofile.reactive.messaging.Message; @ApplicationScoped @@ -28,14 +29,19 @@ public boolean canConvert(Message message, Type type) { @Override public Message convert(Message message, Type type) { Action action = extractAction(message); - Data data; + PreservedData payload; if (Action.UNPUBLISH.equals(action)) { - PreservedData preservedData = dataStore.get(extractKey(message)); - data = preservedData == null ? null : preservedData.getData(); + String key = extractKey(message); + payload = Optional.ofNullable(dataStore.get(key)) + .map(PreservedData::getData) + .map(PreservedData::new) + .orElse(null); } else { - Object payload = message.getPayload(); - data = payload == null ? null : (Data) payload; + payload = Optional.ofNullable(message.getPayload()) + .map(Data.class::cast) + .map(PreservedData::new) + .orElse(null); } - return message.withPayload(new PreservedData(data)); + return message.withPayload(payload); } } diff --git a/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/converter/PreservedRenderingContextMessageConverter.java b/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/converter/PreservedRenderingContextMessageConverter.java index 8c8d6a366..26caa8179 100644 --- a/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/converter/PreservedRenderingContextMessageConverter.java +++ b/rendering-engine-processing-service/src/main/java/dev/streamx/blueprints/rendering/engine/converter/PreservedRenderingContextMessageConverter.java @@ -11,6 +11,7 @@ import io.smallrye.reactive.messaging.MessageConverter; import jakarta.enterprise.context.ApplicationScoped; import java.lang.reflect.Type; +import java.util.Optional; import org.eclipse.microprofile.reactive.messaging.Message; @ApplicationScoped @@ -27,14 +28,19 @@ public boolean canConvert(Message message, Type type) { @Override public Message convert(Message message, Type type) { - RenderingContext renderingContext; + PreservedRenderingContext payload; if (Action.UNPUBLISH.equals(extractAction(message))) { - PreservedRenderingContext entry = renderingContextStore.get(extractKey(message)); - renderingContext = entry == null ? null : entry.getRenderingContext(); + String key = extractKey(message); + payload = Optional.ofNullable(renderingContextStore.get(key)) + .map(PreservedRenderingContext::getRenderingContext) + .map(PreservedRenderingContext::new) + .orElse(null); } else { - Object payload = message.getPayload(); - renderingContext = payload == null ? null : (RenderingContext) payload; + payload = Optional.ofNullable(message.getPayload()) + .map(RenderingContext.class::cast) + .map(PreservedRenderingContext::new) + .orElse(null); } - return message.withPayload(new PreservedRenderingContext(renderingContext)); + return message.withPayload(payload); } } From 17dd237da725bb98270ca754722bd1fbc9fe5aa8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rafa=C5=82=20G=C5=82owacz?= Date: Tue, 21 Oct 2025 12:50:08 +0200 Subject: [PATCH 2/2] DXP-2378 StreamX rendering service does not start with RocksDB enabled * added junit with first unpublish message --- .../engine/RenderingRequestTest.java | 57 +++++++++++++++++-- 1 file changed, 53 insertions(+), 4 deletions(-) diff --git a/rendering-engine-processing-service/src/test/java/dev/streamx/blueprints/rendering/engine/RenderingRequestTest.java b/rendering-engine-processing-service/src/test/java/dev/streamx/blueprints/rendering/engine/RenderingRequestTest.java index b547b8668..70bbde4ae 100644 --- a/rendering-engine-processing-service/src/test/java/dev/streamx/blueprints/rendering/engine/RenderingRequestTest.java +++ b/rendering-engine-processing-service/src/test/java/dev/streamx/blueprints/rendering/engine/RenderingRequestTest.java @@ -8,6 +8,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.wildfly.common.Assert.assertTrue; import dev.streamx.blueprints.data.Data; import dev.streamx.blueprints.data.Renderer; @@ -21,6 +22,7 @@ import io.quarkus.test.junit.QuarkusTest; import io.smallrye.reactive.messaging.memory.InMemorySink; import io.smallrye.reactive.messaging.memory.InMemorySource; +import java.time.Duration; import org.apache.commons.lang3.tuple.Pair; import org.eclipse.microprofile.reactive.messaging.Message; import org.junit.jupiter.api.BeforeEach; @@ -49,12 +51,14 @@ void beforeEach() { @Test void dataPublishRenderingRequestShouldGeneratePage() { String templateKey = "rendering-request-test-page-renderer"; - Pair> data = dataMessage("rendering-request-test-data-type1:1", PUBLISH); + Pair> data = dataMessage( + "rendering-request-test-data-type1:1", PUBLISH); Pair> renderingContextPublishMessage = renderingContextMessage( "rendering-request-test-pages-rendering-context", PUBLISH, - new RenderingContext(templateKey, "rendering-request-test-data-type1:.*", + new RenderingContext(templateKey, + "rendering-request-test-data-type1:.*", null, "rendering-request-test-generated/{{id}}.html", null, @@ -78,15 +82,60 @@ void dataPublishRenderingRequestShouldGeneratePage() { pages.received().get(0)); } + @Test + void dataPublishRenderingRequestShouldGeneratePageAfterInitialUnpublish() { + String templateKey = "rendering-request-test-page-renderer"; + + Pair> renderingContextPublishMessage = + renderingContextMessage( + "rendering-request-test-pages-rendering-context", + PUBLISH, + new RenderingContext(templateKey, + "rendering-request-test-data-type1:.*", + null, + "rendering-request-test-generated/{{id}}.html", + null, + OutputFormat.PAGE)); + + Pair> unpublishData = dataMessage( + "rendering-request-test-data-type1:1", UNPUBLISH); + + renderingContexts.send(renderingContextPublishMessage.getValue()); + renderers.send(rendererPublishMessage(templateKey)); + dataSource.send(unpublishData.getValue()); + await().atLeast(Duration.ofMillis(100)).untilAsserted(() -> + assertTrue(pages.received().isEmpty()) + ); + + Pair> publishData = dataMessage( + "rendering-request-test-data-type1:1", PUBLISH); + + dataSource.send(publishData.getValue()); + sendRenderingRequest( + renderingContextPublishMessage.getKey(), publishData.getKey(), templateKey, PUBLISH, + renderingContextPublishMessage.getValue().getPayload().getOutputKeyTemplate(), + renderingContextPublishMessage.getValue().getPayload().getOutputTypeTemplate()); + + await().until(() -> pages.received().size() == 1); + assertOutput( + "rendering-request-test-generated/" + publishData.getKey() + ".html", + PUBLISH, + null, + "id = " + publishData.getKey(), + pages.received().get(0)); + } + @Test void dataUnpublishRenderingRequestShouldRemovePage() { String templateKey = "rendering-request-test-page-renderer"; - Pair> data = dataMessage("rendering-request-test-data-type2:1", PUBLISH); + Pair> data = dataMessage( + "rendering-request-test-data-type2:1", PUBLISH); Pair> renderingContextPublishMessage = renderingContextMessage( "rendering-request-test-pages-rendering-context", PUBLISH, - new RenderingContext(templateKey, "rendering-request-test-data-type2:.*", + new RenderingContext(templateKey, + "rendering-request-test-data-type2:.*", null, "rendering-request-test-generated/{{id}}.html", "output-template-test-{{id}}",