Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,7 @@ public GenericPayload<Resource> 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());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -85,8 +86,12 @@ public Multi<Message<RenderingRequest>> processRenderer(Message<Renderer> incomi
@Outgoing(Channels.Outgoing.RENDERING_REQUESTS)
public Multi<Message<RenderingRequest>> processContext(
Message<PreservedRenderingContext> 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 {
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -115,12 +118,20 @@ private Stream<Message<RenderingRequest>> calculateRenderingRequestForDataStore(

private Stream<KeyedValue<Data>> 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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -49,12 +51,14 @@ void beforeEach() {
@Test
void dataPublishRenderingRequestShouldGeneratePage() {
String templateKey = "rendering-request-test-page-renderer";
Pair<String, Message<Data>> data = dataMessage("rendering-request-test-data-type1:1", PUBLISH);
Pair<String, Message<Data>> data = dataMessage(
"rendering-request-test-data-type1:1", PUBLISH);
Pair<String, Message<RenderingContext>> 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,
Expand All @@ -78,15 +82,60 @@ void dataPublishRenderingRequestShouldGeneratePage() {
pages.received().get(0));
}

@Test
void dataPublishRenderingRequestShouldGeneratePageAfterInitialUnpublish() {
String templateKey = "rendering-request-test-page-renderer";

Pair<String, Message<RenderingContext>> 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<String, Message<Data>> 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<String, Message<Data>> 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<String, Message<Data>> data = dataMessage("rendering-request-test-data-type2:1", PUBLISH);
Pair<String, Message<Data>> data = dataMessage(
"rendering-request-test-data-type2:1", PUBLISH);
Pair<String, Message<RenderingContext>> 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}}",
Expand Down