diff --git a/README.md b/README.md
index 110e8674..69eb6438 100644
--- a/README.md
+++ b/README.md
@@ -1,8 +1,8 @@
# Code Samples
-This repository contains several Axon-based sample projects to explain specific topics.
-Axon usage extends itself from the repositories under the [Axon Framework project](https://github.com/AxonFramework/),
-as well as the tools provided by AxonIQ like [Axon Server](https://www.axoniq.io/products/axon-server).
+This repository contains several Axon Framework 5-based sample projects to explain specific topics.
+Axon usage extends itself from the repositories under the [Axon Framework project](https://github.com/AxonIQ/AxonFramework),
+as well as the tools provided by AxoniQ like [Axon Server](https://www.axoniq.io/products/axon-server).
Each Maven module within this project represents a different sample you can use, each with its own `README.md`
explaining the intent and usage.
@@ -17,23 +17,31 @@ Down below is an exhaustive list of all the sample:
and Spring-based application.
2. [Distributed Exceptions](distributed-exceptions/README.md) - Sample showing how to deal with exceptions in a
distributed application.
-3. [Multitenancy](multitenancy/README.md) - Sample showing 'a' approach to multitenancy.
-4. [Reset Handler](reset-handler/README.md) - Sample showing how to reset a `StreamingEventProcessor`.
-5. [Saga - No TOAST in PostgreSQL](saga/README.md) - Sample showing a basis Saga, that's stored in a PostgreSQL
- without [TOAST](https://wiki.postgresql.org/wiki/TOAST).
-6. [Serialization Avro](serialization-avro/README.md) - Sample showing usage of Apache Avro Commands/Events/Queries.
-7. [Sequencing Policy](sequencing-policy/README.md) - Sample showing how to set up a custom `SequencingPolicy` to adjust
- the event sequence for a `StreamingEventProcessor`.
-8. [Set-Based Validation](set-based-validation/README.md) - Sample showing several approaches to implement set-based
- validation.
-9. [Set-Based Validation - Actor Model](set-based-validation-actor-model/README.md) - Sample showing how to implement
- set-based validation through a dedicated aggregate instance.
-10. [Snapshots](snapshots/README.md) - Sample showing how to configure aggregate snapshotting.
-11. [Stateful Event Handler](stateful-event-handler/README.md) - Sample showing a stateful event handler that can be
- used as a replacement for sagas.
-12. [Subscription Query - REST](subscription-query-rest/README.md) - Sample showing how to use Axon's subscription query
- cleanly in a REST-based controller.
-13. [Subscription Query - Streaming](subscription-query-streaming/README.md) - Sample showing how to use Axon's
+3. [Order Fulfillment Workflow](order-fulfillment-workflow/README.md) - Sample showing how to model an Order
+ Fulfillment process with the AxoniQ Workflow Engine.
+4. [Reset Handler](reset-handler/README.md) - Sample showing how to reset a `PooledStreamingEventProcessor`.
+5. [Serialization Avro](serialization-avro/README.md) - Sample showing usage of Apache Avro Commands/Events/Queries.
+6. [Sequencing Policy](sequencing-policy/README.md) - Sample showing how to set up a custom `SequencingPolicy` to adjust
+ the event sequence for a `PooledStreamingEventProcessor`.
+7. [Snapshots](snapshots/README.md) - Sample showing how to configure event-sourced entity snapshotting.
+8. [Stateful Event Handler](stateful-event-handler/README.md) - Sample showing a stateful event handler that can be
+ used as a replacement for sagas.
+9. [Subscription Query - REST](subscription-query-rest/README.md) - Sample showing how to use Axon's subscription query
+ cleanly in a REST-based controller.
+10. [Subscription Query - Streaming](subscription-query-streaming/README.md) - Sample showing how to use Axon's
subscription query cleanly in a streaming-based controller.
-14. [Upcasters](upcaster/README.md) - Sample showing several implementations of upcasters.
+11. [Workflow Saga](workflow-saga/README.md) - Sample showing a Saga-like process modelled with the AxoniQ Workflow
+ Engine.
+
+## Topics now covered by the Axon Framework project itself
+
+A few topics that used to live here as standalone samples now have more thorough, actively maintained examples in the
+[AxonIQ/AxonFramework](https://github.com/AxonIQ/AxonFramework) repository's `examples/` directory. Rather than
+maintaining a second, thinner copy, we point to those instead:
+
+- **Sagas** - see the [`extension-workflow`](https://github.com/AxonIQ/extension-workflow) project and the automation
+ patterns in [`examples/university-demo`](https://github.com/AxonIQ/AxonFramework/tree/main/examples/university-demo).
+- **Multitenancy** - see [`examples/university-multi-tenancy-examples`](https://github.com/AxonIQ/AxonFramework/tree/main/examples/university-multi-tenancy-examples).
+- **Set-based validation** - see `CourseUniqueNameSetValidation` in [`examples/university-demo`](https://github.com/AxonIQ/AxonFramework/tree/main/examples/university-demo/src/main/java/org/axonframework/examples/demo/university/faculty/write/createcourse).
+- **Upcasting / event transformation** - see [`examples/university-message-transformation`](https://github.com/AxonIQ/AxonFramework/tree/main/examples/university-message-transformation).
diff --git a/axon-spring-template/pom.xml b/axon-spring-template/pom.xml
index c2c3e222..fd6e2c25 100644
--- a/axon-spring-template/pom.xml
+++ b/axon-spring-template/pom.xml
@@ -14,7 +14,7 @@
- org.axonframework
+ org.axonframework.extensions.springaxon-spring-boot-starter
diff --git a/distributed-exceptions/pom.xml b/distributed-exceptions/pom.xml
index a84c8981..f3714305 100644
--- a/distributed-exceptions/pom.xml
+++ b/distributed-exceptions/pom.xml
@@ -29,7 +29,7 @@
axon-eventsourcing
- org.axonframework
+ org.axonframework.extensions.springaxon-spring-boot-starter
@@ -56,7 +56,7 @@
org.testcontainers
- junit-jupiter
+ testcontainers-junit-jupiterorg.junit.jupiter
diff --git a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/DistributedExceptionsApplication.java b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/DistributedExceptionsApplication.java
index fbb4831b..4a0baada 100644
--- a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/DistributedExceptionsApplication.java
+++ b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/DistributedExceptionsApplication.java
@@ -1,10 +1,7 @@
package io.axoniq.distributedexceptions;
-import io.axoniq.distributedexceptions.command.ExceptionWrappingHandlerInterceptor;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
-import org.springframework.context.annotation.Bean;
-import org.springframework.context.annotation.Profile;
@SpringBootApplication
public class DistributedExceptionsApplication {
@@ -12,10 +9,4 @@ public class DistributedExceptionsApplication {
public static void main(String[] args) {
SpringApplication.run(DistributedExceptionsApplication.class, args);
}
-
- @Bean
- @Profile("command")
- public ExceptionWrappingHandlerInterceptor exceptionWrappingHandlerInterceptor() {
- return new ExceptionWrappingHandlerInterceptor();
- }
}
diff --git a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/CardIssuedEvent.java b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/CardIssuedEvent.java
index c812eac4..eda51bff 100644
--- a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/CardIssuedEvent.java
+++ b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/CardIssuedEvent.java
@@ -1,7 +1,9 @@
package io.axoniq.distributedexceptions.api;
-import javax.annotation.Nonnull;
+import org.axonframework.eventsourcing.annotation.EventTag;
+import org.axonframework.messaging.eventhandling.annotation.Event;
-public record CardIssuedEvent(@Nonnull String id, int amount) {
+@Event
+public record CardIssuedEvent(@EventTag(key = "GiftCard") String id, int amount) {
}
diff --git a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/CardRedeemedEvent.java b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/CardRedeemedEvent.java
index 5c724482..a8c1626d 100644
--- a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/CardRedeemedEvent.java
+++ b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/CardRedeemedEvent.java
@@ -1,7 +1,9 @@
package io.axoniq.distributedexceptions.api;
-import javax.annotation.Nonnull;
+import org.axonframework.eventsourcing.annotation.EventTag;
+import org.axonframework.messaging.eventhandling.annotation.Event;
-public record CardRedeemedEvent(@Nonnull String id, int amount) {
+@Event
+public record CardRedeemedEvent(@EventTag(key = "GiftCard") String id, int amount) {
}
diff --git a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/IssueCardCommand.java b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/IssueCardCommand.java
index 3dcafbfe..c4a34e86 100644
--- a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/IssueCardCommand.java
+++ b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/IssueCardCommand.java
@@ -1,9 +1,9 @@
package io.axoniq.distributedexceptions.api;
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
+import org.axonframework.messaging.commandhandling.annotation.Command;
+import org.axonframework.modelling.annotation.TargetEntityId;
-import javax.annotation.Nonnull;
-
-public record IssueCardCommand(@TargetAggregateIdentifier @Nonnull String id, int amount) {
+@Command(routingKey = "id")
+public record IssueCardCommand(@TargetEntityId String id, int amount) {
}
diff --git a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/RedeemCardCommand.java b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/RedeemCardCommand.java
index 4b369ac1..c36e44aa 100644
--- a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/RedeemCardCommand.java
+++ b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/api/RedeemCardCommand.java
@@ -1,9 +1,9 @@
package io.axoniq.distributedexceptions.api;
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
+import org.axonframework.messaging.commandhandling.annotation.Command;
+import org.axonframework.modelling.annotation.TargetEntityId;
-import javax.annotation.Nonnull;
-
-public record RedeemCardCommand(@TargetAggregateIdentifier @Nonnull String id, int amount) {
+@Command(routingKey = "id")
+public record RedeemCardCommand(@TargetEntityId String id, int amount) {
}
diff --git a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/command/ExceptionWrappingHandlerInterceptor.java b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/command/ExceptionWrappingHandlerInterceptor.java
index 5cf88137..2de71cdd 100644
--- a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/command/ExceptionWrappingHandlerInterceptor.java
+++ b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/command/ExceptionWrappingHandlerInterceptor.java
@@ -2,31 +2,37 @@
import io.axoniq.distributedexceptions.api.GiftCardBusinessError;
import io.axoniq.distributedexceptions.api.GiftCardBusinessErrorCode;
-import org.axonframework.commandhandling.CommandExecutionException;
-import org.axonframework.commandhandling.CommandMessage;
-import org.axonframework.messaging.InterceptorChain;
-import org.axonframework.messaging.MessageHandlerInterceptor;
-import org.axonframework.messaging.unitofwork.UnitOfWork;
+import org.axonframework.messaging.commandhandling.CommandExecutionException;
+import org.axonframework.messaging.commandhandling.CommandMessage;
+import org.axonframework.messaging.core.MessageHandlerInterceptor;
+import org.axonframework.messaging.core.MessageHandlerInterceptorChain;
+import org.axonframework.messaging.core.MessageStream;
+import org.axonframework.messaging.core.unitofwork.ProcessingContext;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import org.springframework.context.annotation.Profile;
+import org.springframework.stereotype.Component;
import java.lang.invoke.MethodHandles;
-import javax.annotation.Nonnull;
-public class ExceptionWrappingHandlerInterceptor implements MessageHandlerInterceptor> {
+@Component
+@Profile("command")
+public class ExceptionWrappingHandlerInterceptor implements MessageHandlerInterceptor {
private static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
@Override
- public Object handle(@Nonnull UnitOfWork extends CommandMessage>> unitOfWork,
- @Nonnull InterceptorChain interceptorChain) {
- try {
- return interceptorChain.proceed();
- } catch (Throwable e) {
- throw new CommandExecutionException(
- "An exception has occurred during command execution", e, exceptionDetails(e)
- );
- }
+ public MessageStream> interceptOnHandle(CommandMessage message,
+ ProcessingContext context,
+ MessageHandlerInterceptorChain chain) {
+ return chain.proceed(message, context)
+ .onErrorContinue(throwable -> MessageStream.failed(
+ new CommandExecutionException(
+ "An exception has occurred during command execution",
+ throwable,
+ exceptionDetails(throwable)
+ )
+ ));
}
// Domain specific details can be returned in a couple of forms.
diff --git a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/command/GiftCard.java b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/command/GiftCard.java
index 4b27cccc..3a9b283c 100644
--- a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/command/GiftCard.java
+++ b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/command/GiftCard.java
@@ -4,39 +4,37 @@
import io.axoniq.distributedexceptions.api.CardIssuedEvent;
import io.axoniq.distributedexceptions.api.RedeemCardCommand;
import io.axoniq.distributedexceptions.api.CardRedeemedEvent;
-import org.axonframework.commandhandling.CommandHandler;
-import org.axonframework.eventsourcing.EventSourcingHandler;
-import org.axonframework.modelling.command.AggregateIdentifier;
-import org.axonframework.spring.stereotype.Aggregate;
+import org.axonframework.eventsourcing.annotation.EventSourcingHandler;
+import org.axonframework.eventsourcing.annotation.reflection.EntityCreator;
+import org.axonframework.extension.spring.stereotype.EventSourced;
+import org.axonframework.messaging.commandhandling.annotation.CommandHandler;
+import org.axonframework.messaging.eventhandling.gateway.EventAppender;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Profile;
import java.lang.invoke.MethodHandles;
-import static org.axonframework.modelling.command.AggregateLifecycle.apply;
-
-@Aggregate
+@EventSourced(tagKey = "GiftCard")
@Profile("command")
class GiftCard {
private static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
- @AggregateIdentifier
private String giftCardId;
private int remainingValue;
@CommandHandler
- public GiftCard(IssueCardCommand command) {
+ public static void handle(IssueCardCommand command, EventAppender eventAppender) {
logger.debug("handling {}", command);
if (command.amount() <= 0) {
throw new NegativeOrZeroAmount(command.amount(), "amount <= 0");
}
- apply(new CardIssuedEvent(command.id(), command.amount()));
+ eventAppender.append(new CardIssuedEvent(command.id(), command.amount()));
}
@CommandHandler
- public void handle(RedeemCardCommand command) {
+ public void handle(RedeemCardCommand command, EventAppender eventAppender) {
logger.debug("handling {}", command);
if (command.amount() <= 0) {
throw new NegativeOrZeroAmount(command.amount(), "amount <= 0");
@@ -44,7 +42,7 @@ public void handle(RedeemCardCommand command) {
if (command.amount() > remainingValue) {
throw new InsufficientFunds("amount > remaining value");
}
- apply(new CardRedeemedEvent(giftCardId, command.amount()));
+ eventAppender.append(new CardRedeemedEvent(giftCardId, command.amount()));
}
@EventSourcingHandler
@@ -62,6 +60,7 @@ public void on(CardRedeemedEvent event) {
logger.debug("new remaining value: {}", remainingValue);
}
+ @EntityCreator
public GiftCard() {
// Required by Axon
logger.debug("Empty constructor invoked");
diff --git a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/rest/GiftCardController.java b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/rest/GiftCardController.java
index 77916dd2..77832c79 100644
--- a/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/rest/GiftCardController.java
+++ b/distributed-exceptions/src/main/java/io/axoniq/distributedexceptions/rest/GiftCardController.java
@@ -3,8 +3,8 @@
import io.axoniq.distributedexceptions.api.GiftCardBusinessError;
import io.axoniq.distributedexceptions.api.IssueCardCommand;
import io.axoniq.distributedexceptions.api.RedeemCardCommand;
-import org.axonframework.commandhandling.CommandExecutionException;
-import org.axonframework.commandhandling.gateway.CommandGateway;
+import org.axonframework.messaging.commandhandling.CommandExecutionException;
+import org.axonframework.messaging.commandhandling.gateway.CommandGateway;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Profile;
@@ -37,7 +37,8 @@ public GiftCardController(CommandGateway commandGateway) {
public CompletableFuture> issueNewGiftCard(@RequestBody IssueCardDto request) {
IssueCardCommand command = new IssueCardCommand(UUID.randomUUID().toString(), request.amount());
return commandGateway.send(command)
- .thenApply(it -> ResponseEntity.ok(String.valueOf(it)))
+ .getResultMessage()
+ .thenApply(it -> ResponseEntity.ok(String.valueOf(it.payload())))
.exceptionally(e -> {
logException(e);
String errorResponse = getErrorResponseMessage(e.getCause());
@@ -49,6 +50,7 @@ public CompletableFuture> issueNewGiftCard(@RequestBody I
@PutMapping("/{id}")
public CompletableFuture> redeem(@PathVariable String id, @RequestBody RedeemCardDto dto) {
return commandGateway.send(new RedeemCardCommand(id, dto.amount()))
+ .getResultMessage()
.thenApply(it -> ResponseEntity.ok(""))
.exceptionally(e -> {
logException(e);
@@ -68,9 +70,8 @@ private void logException(Throwable throwable) {
private String getErrorResponseMessage(Throwable throwable) {
if (throwable instanceof CommandExecutionException cee) {
- return cee.getDetails()
- .map((Object it) -> {
- GiftCardBusinessError giftCardBusinessError = (GiftCardBusinessError) it;
+ return cee.getDetails(GiftCardBusinessError.class)
+ .map(giftCardBusinessError -> {
logger.debug("Received BusinessError with data: " + giftCardBusinessError);
logger.error("Unable to create GiftCard due to validation constrains. Reason: "
+ giftCardBusinessError);
diff --git a/distributed-exceptions/src/main/resources/application.properties b/distributed-exceptions/src/main/resources/application.properties
index c662bc56..f91296e5 100644
--- a/distributed-exceptions/src/main/resources/application.properties
+++ b/distributed-exceptions/src/main/resources/application.properties
@@ -1,2 +1,2 @@
spring.main.web-application-type=none
-axon.serializer.general=jackson
+axon.converter.general=jackson
diff --git a/multitenancy/.gitignore b/multitenancy/.gitignore
deleted file mode 100644
index 9f5ab22c..00000000
--- a/multitenancy/.gitignore
+++ /dev/null
@@ -1 +0,0 @@
-src/main/frontend
diff --git a/multitenancy/README.md b/multitenancy/README.md
deleted file mode 100644
index 55c1384f..00000000
--- a/multitenancy/README.md
+++ /dev/null
@@ -1,54 +0,0 @@
-# Multi-tenant playground
-
-This simple app allows you to experiment with multi-tenancy feature in Axon Framework.
-By using [Multitenancy extension](https://github.com/AxonFramework/extension-multitenancy) this app can connect to multiple contexts at once dynamicly.
-Using Vaadin UI, explore how to create tenants in runtime and dispatch messages to them.
-
-## Requirements
-
- - Requires Axon Server 2024.1+ with a valid license
- - Axon Framework 4.10.1+
- - [Multitenancy extension](https://github.com/AxonFramework/extension-multitenancy)
- - Uses In-Memory H2 Database for projections - DB per tenant
-
-## How to start
-
-While in `src/main/docker` type in terminal `docker-compose up` to start Axon Server EE.
-
-Make sure to replace content of `axoniq.license` with a valid license.
-
-Start application with `mvn spring-boot:run` or manually via favorite IDE.
-
-After starting application visit and interact with UI at [localhost:8080](http://localhost:8080)
-
-You may start by registering a new tenant.
-
-### With UI you can:
- - Login as existing tenant
-
- - Register new tenant at runtime
- 
- - Send commands, queries, events
- 
- - Reset projections
-
- - Explore H2 Database used for projections - per tenant
- - 
-
-## How it works
-
-### Multi-tenant configuration
-
-[MultiTenantConfig#tenantFilter](src/main/java/io/axoniq/multitenancy/MultiTenantConfig.java) configures application to connect to any context which name starts with `tenant-`.
-If new context is created during runtime that matches the filter, application will automatically connect to that tenant.
-
-[MultiTenantConfig#tenantDataSourceResolver](src/main/java/io/axoniq/multitenancy/MultiTenantConfig.java) configures application to use DB for projections per tenant. It defines datasource properties for tenant.
-Important: DB needs to be created and schema needs to be initialiesd before context is created, otherwise you will see exception such as DB/Table not found...
-
-[CreateTenantService#createTenant](src/main/java/io/axoniq/multitenancy/web/CreateTenantService.java) uses Admin API to interact with Axon Server and create tenant context.
-Before creating context it creates embedded H2 database. As H2 database is in memory, it will not be persistent after application restart, while Axon Server contexts will.
-Be sure to clean up contexts if you are planing to play with this application multiple times or you may expect to see some exceptions.
-
-[MessageService](src/main/java/io/axoniq/multitenancy/web/MessageService.java) routes messages to specific tenant by setting `TENANT_CORRELATION_KEY` MetaData to initial message.
-
-[ResetService](src/main/java/io/axoniq/multitenancy/web/ResetService.java) lists all event processors and resets one for specificity tenant. Convention is that Event Processor name contains original event processor name and tenant name. Such as `even_processor_name@tenant_name`.
diff --git a/multitenancy/actions.png b/multitenancy/actions.png
deleted file mode 100644
index e07959cc..00000000
Binary files a/multitenancy/actions.png and /dev/null differ
diff --git a/multitenancy/h2.png b/multitenancy/h2.png
deleted file mode 100644
index 54fae173..00000000
Binary files a/multitenancy/h2.png and /dev/null differ
diff --git a/multitenancy/login.png b/multitenancy/login.png
deleted file mode 100644
index db2a8d8e..00000000
Binary files a/multitenancy/login.png and /dev/null differ
diff --git a/multitenancy/pom.xml b/multitenancy/pom.xml
deleted file mode 100644
index 5c52c6c0..00000000
--- a/multitenancy/pom.xml
+++ /dev/null
@@ -1,161 +0,0 @@
-
-
- 4.0.0
-
- code-samples
- io.axoniq
- 0.0.2-SNAPSHOT
-
-
- multitenancy
-
- Multitenancy
- Axon Framework application using the Multitenancy Extension
-
-
-
- 11.13.2
- 2.4.240
-
- 24.9.2
- 21
-
-
-
-
-
- com.vaadin
- vaadin-bom
-
- ${vaadin.version}
- pom
- import
-
-
-
-
-
-
-
- org.axonframework
- axon-spring-boot-starter
-
-
- org.axonframework
- axon-configuration
-
-
- io.axoniq
- axonserver-connector-java
-
-
- org.axonframework
- axon-test
-
-
- org.axonframework.extensions.multitenancy
- axon-multitenancy-spring-boot-starter
-
-
- org.axonframework.extensions.multitenancy
- axon-multitenancy-spring-boot-autoconfigure
-
-
- org.axonframework.extensions.multitenancy
- axon-multitenancy
-
-
-
-
- org.springframework.data
- spring-data-commons
-
-
- org.springframework.data
- spring-data-jpa
-
-
- org.springframework.boot
- spring-boot-starter-data-jpa
-
-
- org.springframework.boot
- spring-boot-starter-test
- test
-
-
- org.springframework.boot
- spring-boot-starter-web
-
-
-
- org.flywaydb
- flyway-core
- ${flyway.version}
-
-
- com.h2database
- h2
- ${h2.version}
-
-
-
- io.projectreactor
- reactor-core
-
-
- com.vaadin
- vaadin-spring-boot-starter
-
-
-
-
-
-
-
- org.apache.maven.plugins
- maven-compiler-plugin
-
- 17
- 17
- 21
-
-
-
-
- com.vaadin
- vaadin-maven-plugin
- ${vaadin.version}
-
-
-
- prepare-frontend
- build-frontend
-
-
-
-
-
-
-
-
-
- org.springframework.boot
- spring-boot-maven-plugin
-
-
-
- repackage
-
-
- spring-boot
- io.axoniq.multitenancy.MultitenancyExampleApplication
-
-
-
-
-
-
-
diff --git a/multitenancy/register.png b/multitenancy/register.png
deleted file mode 100644
index 7a180ba2..00000000
Binary files a/multitenancy/register.png and /dev/null differ
diff --git a/multitenancy/reset.png b/multitenancy/reset.png
deleted file mode 100644
index f96efe3d..00000000
Binary files a/multitenancy/reset.png and /dev/null differ
diff --git a/multitenancy/src/main/docker/axoniq.license b/multitenancy/src/main/docker/axoniq.license
deleted file mode 100644
index 03a560a7..00000000
--- a/multitenancy/src/main/docker/axoniq.license
+++ /dev/null
@@ -1,14 +0,0 @@
-#AxonIQ License File generated on 2022-09-20
-#Tue Sep 20 08:26:45 UTC 2022
-license_key_id=dd722d00-dd7c-4c70-bf1a-2f5c78a3af6c
-signature=POSseYXhv/Op+DaXT7UYjQudjzhkM3adU/gPGWuSEUPXTspHoF788iL0poZoycIVBsm8CGxUqeWdZyAaWQnvEuveqQCl9YlgmK6ocYp6uPL3ltxfAO93xW8X9paUqMMj1Y6FFky8ITSPRPWNA1zpTFyoPLs3Dmy+/kx/DjVJBiTo8QaDt08IWhr7f8Zezp9n6HgZKU/UnojSwBBtp/NxxMGu11VjjqrBI6NBBj+J3LtJMNkA5P8za4cQ2Huu+HufROz+6N2Ia1obdkpG2eqaoRsWf6MLKzYBX6eU9nV/b5SFCOX9/ehw3O0lHT5WZokY3ZKZUaXTAWNYgTopRkpeunNv7p+pVHnSAV969zWZ6Y1YbgDkL1TAmnqWu1+1w3C48xoyFvlSSNcUkhWljOlwF2X0gmYX6/r02UDPD5Q0nQcqC01QcOFOeYny0VgWe6XfYovKMmlqDCW5m2y/W7dEJffpl4I93p0OFzh4LRxqrqriNTVKs95YNh5XPvm4e8u2/0ieeo9so8gCz+xTeDB6rZw81YCyLDYkS4GP1pVt/SOPtvZHPncKWOdv19790zc6azByyXAxnqi/GXGI6dB9kaKDmBVPeCMtSdF6Z0Cq9K+aGAHoFhwZZJmG4Bo1MAutAAA6J+Qw60yilehKt+SgibwQFxIn1LEkRZu4Ga4e/mg\=
-reference=Trial license for use by the participants to AxonIQCon 2022
-contexts=100
-expiry_date=2022-10-15
-contact=bert.laverman@axoniq.io
-edition=Enterprise
-issue_date=2022-09-20
-grace_date=2022-10-15
-product=AxonServer
-licensee=AxonIQCon 2022
-clusterNodes=3
\ No newline at end of file
diff --git a/multitenancy/src/main/docker/cluster-template.yml b/multitenancy/src/main/docker/cluster-template.yml
deleted file mode 100644
index 19c609e1..00000000
--- a/multitenancy/src/main/docker/cluster-template.yml
+++ /dev/null
@@ -1,33 +0,0 @@
-axoniq:
- axonserver:
- cluster-template:
- users: []
- replicationGroups:
- - roles:
- - role: PRIMARY
- node: axonserver-enterprise-1
- - role: PRIMARY
- node: axonserver-enterprise-2
- - role: PRIMARY
- node: axonserver-enterprise-3
- name: _admin
- contexts:
- - name: _admin
- metaData:
- event.index-format: JUMP_SKIP_INDEX
- snapshot.index-format: JUMP_SKIP_INDEX
- - roles:
- - role: PRIMARY
- node: axonserver-enterprise-2
- - role: PRIMARY
- node: axonserver-enterprise-1
- - role: PRIMARY
- node: axonserver-enterprise-3
- name: default
- contexts:
- - name: tenant-one
- metaData:
- event.index-format: JUMP_SKIP_INDEX
- snapshot.index-format: JUMP_SKIP_INDEX
- applications: []
- first: axonserver-enterprise-1:8224
diff --git a/multitenancy/src/main/docker/docker-compose.yml b/multitenancy/src/main/docker/docker-compose.yml
deleted file mode 100644
index 69d5e01e..00000000
--- a/multitenancy/src/main/docker/docker-compose.yml
+++ /dev/null
@@ -1,84 +0,0 @@
-version: '3.8'
-services:
- axonserver-enterprise-1:
- image: axoniq/axonserver-enterprise:4.6.4-dev
- hostname: axonserver-enterprise-1
- environment:
- - AXONIQ_AXONSERVER_CLUSTERTEMPLATE_PATH=/axonserver/cluster-template.yml
- - AXONIQ_AXONSERVER_ENTERPRISE_LICENSE-DIRECTORY=/axonserver
- - SERVER_PORT=8024
- - AXONIQ_AXONSERVER_PORT=8124
- - AXONIQ_AXONSERVER_METRICS_GRPC_ENABLED=true
- - AXONIQ_AXONSERVER_METRICS_GRPC_PROMETHEUS-ENABLED=true
- volumes:
- - axonserver-enterprise-1-log:/axonserver/log
- - axonserver-enterprise-1-events:/axonserver/events
- - axonserver-enterprise-1-data:/axonserver/data
- - ./axoniq.license:/axonserver/axoniq.license
- - ./cluster-template.yml:/axonserver/cluster-template.yml
- ports:
- - '8024:8024'
- - '8124:8124'
- - '8224:8224'
- networks:
- - axon-net
-
- axonserver-enterprise-2:
- image: axoniq/axonserver-enterprise:4.6.4-dev
- hostname: axonserver-enterprise-2
- environment:
- - AXONIQ_AXONSERVER_CLUSTERTEMPLATE_PATH=/axonserver/cluster-template.yml
- - AXONIQ_AXONSERVER_ENTERPRISE_LICENSE-DIRECTORY=/axonserver
- - SERVER_PORT=8025
- - AXONIQ_AXONSERVER_PORT=8125
- - AXONIQ_AXONSERVER_METRICS_GRPC_ENABLED=true
- - AXONIQ_AXONSERVER_METRICS_GRPC_PROMETHEUS-ENABLED=true
- volumes:
- - axonserver-enterprise-2-log:/axonserver/log
- - axonserver-enterprise-2-events:/axonserver/events
- - axonserver-enterprise-2-data:/axonserver/data
- - ./axoniq.license:/axonserver/axoniq.license
- - ./cluster-template.yml:/axonserver/cluster-template.yml
- ports:
- - '8025:8025'
- - '8125:8125'
- - '8225:8225'
- networks:
- - axon-net
-
- axonserver-enterprise-3:
- image: axoniq/axonserver-enterprise:4.6.4-dev
- hostname: axonserver-enterprise-3
- environment:
- - AXONIQ_AXONSERVER_CLUSTERTEMPLATE_PATH=/axonserver/cluster-template.yml
- - AXONIQ_AXONSERVER_ENTERPRISE_LICENSE-DIRECTORY=/axonserver
- - SERVER_PORT=8026
- - AXONIQ_AXONSERVER_PORT=8126
- - AXONIQ_AXONSERVER_METRICS_GRPC_ENABLED=true
- - AXONIQ_AXONSERVER_METRICS_GRPC_PROMETHEUS-ENABLED=true
- volumes:
- - axonserver-enterprise-3-log:/axonserver/log
- - axonserver-enterprise-3-events:/axonserver/events
- - axonserver-enterprise-3-data:/axonserver/data
- - ./axoniq.license:/axonserver/axoniq.license
- - ./cluster-template.yml:/axonserver/cluster-template.yml
- ports:
- - '8026:8026'
- - '8126:8126'
- - '8226:8226'
- networks:
- - axon-net
-
-volumes:
- axonserver-enterprise-1-log:
- axonserver-enterprise-1-events:
- axonserver-enterprise-1-data:
- axonserver-enterprise-2-log:
- axonserver-enterprise-2-events:
- axonserver-enterprise-2-data:
- axonserver-enterprise-3-log:
- axonserver-enterprise-3-events:
- axonserver-enterprise-3-data:
-
-networks:
- axon-net:
\ No newline at end of file
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/MultiTenantConfig.java b/multitenancy/src/main/java/io/axoniq/multitenancy/MultiTenantConfig.java
deleted file mode 100644
index a539af43..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/MultiTenantConfig.java
+++ /dev/null
@@ -1,43 +0,0 @@
-package io.axoniq.multitenancy;
-
-import org.axonframework.axonserver.connector.AxonServerConfiguration;
-import org.axonframework.axonserver.connector.AxonServerConnectionManager;
-import org.axonframework.axonserver.connector.event.axon.AxonServerEventStore;
-import org.axonframework.eventsourcing.snapshotting.SnapshotFilter;
-import org.axonframework.extensions.multitenancy.components.TenantConnectPredicate;
-
-import org.axonframework.extensions.multitenancy.configuration.MultiTenantStreamableMessageSourceProvider;
-import org.springframework.context.annotation.Bean;
-import org.springframework.context.annotation.Configuration;
-
-import java.util.concurrent.Executors;
-import java.util.concurrent.ScheduledExecutorService;
-
-@Configuration
-public class MultiTenantConfig {
-
- @Bean
- public TenantConnectPredicate tenantFilter() {
- return context -> context.tenantId().startsWith("tenant-");
- }
-
- @Bean
- public ScheduledExecutorService persistentStreamScheduler() {
- return Executors.newScheduledThreadPool(10, Thread.ofVirtual()
- .name("persistent-streams-", 0)
- .factory());
- }
-
- // UNCOMMENT THIS BEAN TO ENABLE MULTITENANCY WITH MULTIPLE DATA SOURCES
-// @Bean
-// public Function tenantDataSourceResolver() {
-// return tenant -> {
-// DataSourceProperties properties = new DataSourceProperties();
-// properties.setUrl("jdbc:h2:mem:" + tenant.tenantId());
-// properties.setDriverClassName("org.h2.Driver");
-// properties.setUsername("sa");
-// return properties;
-// };
-// }
-}
-
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/MultitenancyExampleApplication.java b/multitenancy/src/main/java/io/axoniq/multitenancy/MultitenancyExampleApplication.java
deleted file mode 100644
index 3b7113f7..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/MultitenancyExampleApplication.java
+++ /dev/null
@@ -1,14 +0,0 @@
-package io.axoniq.multitenancy;
-
-import org.springframework.boot.SpringApplication;
-import org.springframework.boot.autoconfigure.SpringBootApplication;
-import org.springframework.data.jpa.repository.config.EnableJpaRepositories;
-
-@SpringBootApplication
-@EnableJpaRepositories
-public class MultitenancyExampleApplication {
-
- public static void main(String[] args) {
- SpringApplication.run(MultitenancyExampleApplication.class, args);
- }
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/api/AddFundsCommand.java b/multitenancy/src/main/java/io/axoniq/multitenancy/api/AddFundsCommand.java
deleted file mode 100644
index 6cb2e001..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/api/AddFundsCommand.java
+++ /dev/null
@@ -1,9 +0,0 @@
-package io.axoniq.multitenancy.api;
-
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
-
-import javax.annotation.Nonnull;
-
-public record AddFundsCommand(@TargetAggregateIdentifier @Nonnull String id, Integer amount) {
-
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/api/CardIssuedEvent.java b/multitenancy/src/main/java/io/axoniq/multitenancy/api/CardIssuedEvent.java
deleted file mode 100644
index 4c488ab3..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/api/CardIssuedEvent.java
+++ /dev/null
@@ -1,5 +0,0 @@
-package io.axoniq.multitenancy.api;
-
-public record CardIssuedEvent(String id, Integer amount) {
-
-}
\ No newline at end of file
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/api/FindAllCardsQuery.java b/multitenancy/src/main/java/io/axoniq/multitenancy/api/FindAllCardsQuery.java
deleted file mode 100644
index 57219514..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/api/FindAllCardsQuery.java
+++ /dev/null
@@ -1,5 +0,0 @@
-package io.axoniq.multitenancy.api;
-
-public record FindAllCardsQuery() {
-
-}
\ No newline at end of file
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/api/FindCardQuery.java b/multitenancy/src/main/java/io/axoniq/multitenancy/api/FindCardQuery.java
deleted file mode 100644
index 669a58c1..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/api/FindCardQuery.java
+++ /dev/null
@@ -1,5 +0,0 @@
-package io.axoniq.multitenancy.api;
-
-public record FindCardQuery(String id) {
-
-}
\ No newline at end of file
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/api/FundsAddedEvent.java b/multitenancy/src/main/java/io/axoniq/multitenancy/api/FundsAddedEvent.java
deleted file mode 100644
index 31f815bf..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/api/FundsAddedEvent.java
+++ /dev/null
@@ -1,7 +0,0 @@
-package io.axoniq.multitenancy.api;
-
-import javax.annotation.Nonnull;
-
-public record FundsAddedEvent(@Nonnull String id, int amount) {
-
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/api/GiftCardRecord.java b/multitenancy/src/main/java/io/axoniq/multitenancy/api/GiftCardRecord.java
deleted file mode 100644
index 7e75db64..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/api/GiftCardRecord.java
+++ /dev/null
@@ -1,5 +0,0 @@
-package io.axoniq.multitenancy.api;
-
-public record GiftCardRecord(String id, Integer initialValue, Integer remainingValue, String payload) {
-
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/api/IssueCardCommand.java b/multitenancy/src/main/java/io/axoniq/multitenancy/api/IssueCardCommand.java
deleted file mode 100644
index a1270c7e..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/api/IssueCardCommand.java
+++ /dev/null
@@ -1,7 +0,0 @@
-package io.axoniq.multitenancy.api;
-
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
-
-public record IssueCardCommand(@TargetAggregateIdentifier String id, Integer amount) {
-
-}
\ No newline at end of file
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/domain/GiftCard.java b/multitenancy/src/main/java/io/axoniq/multitenancy/domain/GiftCard.java
deleted file mode 100644
index 9513caf3..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/domain/GiftCard.java
+++ /dev/null
@@ -1,64 +0,0 @@
-package io.axoniq.multitenancy.domain;
-
-import io.axoniq.multitenancy.api.AddFundsCommand;
-import io.axoniq.multitenancy.api.FundsAddedEvent;
-import io.axoniq.multitenancy.api.IssueCardCommand;
-import io.axoniq.multitenancy.api.CardIssuedEvent;
-import org.axonframework.commandhandling.CommandHandler;
-import org.axonframework.config.ProcessingGroup;
-import org.axonframework.eventsourcing.EventSourcingHandler;
-import org.axonframework.modelling.command.AggregateIdentifier;
-import org.axonframework.spring.stereotype.Aggregate;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.lang.invoke.MethodHandles;
-
-import static org.axonframework.modelling.command.AggregateLifecycle.apply;
-
-@ProcessingGroup("giftcard")
-@Aggregate
-class GiftCard {
-
- private final static Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
-
- @AggregateIdentifier
- private String id;
- private int remainingValue;
-
- public GiftCard() {
- logger.debug("Empty constructor invoked");
- }
-
- @CommandHandler
- GiftCard(IssueCardCommand command) throws InterruptedException {
- logger.debug("handling {}", command);
- Integer amount = command.amount();
- if (amount <= 0) {
- throw new IllegalArgumentException("amount <= 0");
- }
-
- //simulate some work time
- Thread.sleep(amount);
- apply(new CardIssuedEvent(command.id(), amount));
- }
-
- @CommandHandler
- void add(AddFundsCommand command) {
- apply(new FundsAddedEvent(command.id(), command.amount()));
- }
-
- @EventSourcingHandler
- void on(FundsAddedEvent event) {
- remainingValue += event.amount();
- logger.info("remaining amount is {} for account {}", remainingValue, id);
- }
-
- @EventSourcingHandler
- void on(CardIssuedEvent event) {
- logger.debug("applying {}", event);
- id = event.id();
- remainingValue = event.amount();
- logger.debug("new remaining value: {}", remainingValue);
- }
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/query/GiftCardEntity.java b/multitenancy/src/main/java/io/axoniq/multitenancy/query/GiftCardEntity.java
deleted file mode 100644
index fe68393a..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/query/GiftCardEntity.java
+++ /dev/null
@@ -1,34 +0,0 @@
-package io.axoniq.multitenancy.query;
-
-import jakarta.persistence.Entity;
-import jakarta.persistence.Id;
-
-@Entity
-class GiftCardEntity {
-
- @Id
- private String id;
- private Integer initialValue;
- private Integer remainingValue;
-
- public GiftCardEntity(String id, Integer initialValue, Integer remainingValue) {
- this.id = id;
- this.initialValue = initialValue;
- this.remainingValue = remainingValue;
- }
-
- public GiftCardEntity() {
- }
-
- public String getId() {
- return id;
- }
-
- public Integer getInitialValue() {
- return initialValue;
- }
-
- public Integer getRemainingValue() {
- return remainingValue;
- }
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/query/GiftCardHandler.java b/multitenancy/src/main/java/io/axoniq/multitenancy/query/GiftCardHandler.java
deleted file mode 100644
index b60f3803..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/query/GiftCardHandler.java
+++ /dev/null
@@ -1,59 +0,0 @@
-package io.axoniq.multitenancy.query;
-
-import io.axoniq.multitenancy.api.FindAllCardsQuery;
-import io.axoniq.multitenancy.api.FindCardQuery;
-import io.axoniq.multitenancy.api.GiftCardRecord;
-import io.axoniq.multitenancy.api.CardIssuedEvent;
-import org.axonframework.config.ProcessingGroup;
-import org.axonframework.eventhandling.EventHandler;
-import org.axonframework.messaging.MetaData;
-import org.axonframework.queryhandling.QueryHandler;
-import org.axonframework.queryhandling.QueryUpdateEmitter;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import org.springframework.stereotype.Component;
-
-import java.lang.invoke.MethodHandles;
-import java.util.Collections;
-import java.util.List;
-import java.util.Objects;
-import java.util.Optional;
-
-@ProcessingGroup("giftcard")
-@Component
-class GiftCardHandler {
-
- private static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
-
- private final GiftCardJpaRepository giftCardJpaRepository;
-
- public GiftCardHandler(GiftCardJpaRepository giftCardJpaRepository) {
- this.giftCardJpaRepository = giftCardJpaRepository;
- }
-
- @EventHandler
- public void on(CardIssuedEvent event, QueryUpdateEmitter queryUpdateEmitter) {
- /*
- * Update our read model by inserting the new card.
- */
- giftCardJpaRepository.save(new GiftCardEntity(event.id(), event.amount(), event.amount()));
-
- /* Send it to subscription queries of type FindGiftCardQry, but only if the card id matches. */
- queryUpdateEmitter.emit(FindCardQuery.class,
- query -> Objects.equals(event.id(), query.id()),
- new GiftCardRecord(event.id(), event.amount(), event.amount(), "payload")
- );
- }
-
- @QueryHandler
- public List handle(FindAllCardsQuery query, MetaData metaData) {
- logger.debug("@" + metaData + " - " + "FindGiftCardQry: " + query);
- return Collections.emptyList();
- }
-
- @QueryHandler
- public Optional handle(FindCardQuery query, MetaData metaData) {
- logger.debug("@" + metaData + " - " + "FindGiftCardQry: " + query);
- return Optional.empty();
- }
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/query/GiftCardJpaRepository.java b/multitenancy/src/main/java/io/axoniq/multitenancy/query/GiftCardJpaRepository.java
deleted file mode 100644
index f31115b7..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/query/GiftCardJpaRepository.java
+++ /dev/null
@@ -1,9 +0,0 @@
-package io.axoniq.multitenancy.query;
-
-import org.springframework.data.repository.CrudRepository;
-import org.springframework.stereotype.Repository;
-
-@Repository
-public interface GiftCardJpaRepository extends CrudRepository {
-
-}
\ No newline at end of file
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/ui/ActionView.java b/multitenancy/src/main/java/io/axoniq/multitenancy/ui/ActionView.java
deleted file mode 100644
index ee7973f4..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/ui/ActionView.java
+++ /dev/null
@@ -1,129 +0,0 @@
-package io.axoniq.multitenancy.ui;
-
-import com.vaadin.flow.component.Text;
-import com.vaadin.flow.component.UI;
-import com.vaadin.flow.component.button.Button;
-import com.vaadin.flow.component.button.ButtonVariant;
-import com.vaadin.flow.component.details.Details;
-import com.vaadin.flow.component.dialog.Dialog;
-import com.vaadin.flow.component.html.H1;
-import com.vaadin.flow.component.html.Span;
-import com.vaadin.flow.component.notification.Notification;
-import com.vaadin.flow.component.notification.NotificationVariant;
-import com.vaadin.flow.component.orderedlayout.VerticalLayout;
-import com.vaadin.flow.router.BeforeEvent;
-import com.vaadin.flow.router.HasUrlParameter;
-import com.vaadin.flow.router.Route;
-import io.axoniq.multitenancy.web.MessageService;
-import io.axoniq.multitenancy.web.ResetService;
-
-@Route("action")
-public class ActionView extends VerticalLayout implements HasUrlParameter {
-
- private final MessageService messageService;
- private final ResetService resetService;
-
- public ActionView(MessageService messageService,
- ResetService resetService) {
- this.messageService = messageService;
- this.resetService = resetService;
- }
-
- @Override
- public void setParameter(BeforeEvent beforeEvent, String tenantName) {
- VerticalLayout layout = new VerticalLayout();
-
- H1 title = new H1("Actions (" + tenantName + ")");
- layout.add(title);
-
- layout.add(new Text("Do something:"));
-
- layout.add(new Button("Send commands...", evt -> {
- messageService.sendCommands(tenantName);
- showSuccess(tenantName);
- }));
-
-
- layout.add(new Button("Publish events...", evt -> {
- messageService.sendEvents(tenantName);
- showSuccess(tenantName);
- }));
-
- layout.add(new Button("Schedule events...", evt -> {
- messageService.scheduleEvents(tenantName);
- showSuccess(tenantName);
- }));
-
- layout.add(new Button("Send queries...", evt -> {
- messageService.sendQueries(tenantName);
- showSuccess(tenantName);
- }));
-
- layout.add(new Button("Send subscription query...", evt -> {
- messageService.subscriptionQuery(tenantName).block();
- showSuccess(tenantName);
- }));
-
- layout.add(new Button("Reset projections...", evt -> {
- Dialog dialog = new Dialog();
-
- dialog.setHeaderTitle("Reset");
- VerticalLayout content = new VerticalLayout();
- content.setSpacing(false);
- content.setPadding(false);
-
- resetService.listTenantEventProcessors(tenantName)
- .forEach(ep -> {
- Button resetEp = new Button(ep, e -> {
- resetService.reset(ep);
- showSuccess(tenantName);
- });
- content.add(resetEp);
- });
-
- dialog.add(content);
- Button cancelButton = new Button("Close", e -> dialog.close());
- dialog.getFooter().add(cancelButton);
- dialog.open();
- }));
-
- layout.add(new Button("Explore projections...", evt -> {
- Dialog dialog = new Dialog();
-
- dialog.setHeaderTitle("Explore projections");
-
- Span name = new Span("JDBC URL: jdbc:h2:mem:" + tenantName);
- Span email = new Span("username: sa");
-
- VerticalLayout content = new VerticalLayout(name, email);
- content.setSpacing(false);
- content.setPadding(false);
-
- Details details = new Details("Login details", content);
- details.setOpened(true);
- dialog.add(content);
- Button goButton = new Button("Open DB",
- e -> UI.getCurrent().getPage()
- .open("http://localhost:8080/h2-console", "_blank"));
- Button cancelButton = new Button("Close", e -> dialog.close());
- dialog.getFooter().add(goButton);
- dialog.getFooter().add(cancelButton);
- dialog.open();
- }));
-
- layout.setAlignItems(Alignment.CENTER);
-
- Button register = new Button("Or go back...", event -> UI.getCurrent().navigate(""));
- register.addThemeVariants(ButtonVariant.LUMO_CONTRAST);
- layout.add(register);
-
- add(layout);
- }
-
- private void showSuccess(String tenantName) {
- Notification notification = Notification.show(
- "Action done! Check Axon Server dashboard and stats/events for context: " + tenantName
- );
- notification.addThemeVariants(NotificationVariant.LUMO_SUCCESS);
- }
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/ui/MainView.java b/multitenancy/src/main/java/io/axoniq/multitenancy/ui/MainView.java
deleted file mode 100644
index 1b9f6e9a..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/ui/MainView.java
+++ /dev/null
@@ -1,40 +0,0 @@
-package io.axoniq.multitenancy.ui;
-
-import com.vaadin.flow.component.Text;
-import com.vaadin.flow.component.UI;
-import com.vaadin.flow.component.button.Button;
-import com.vaadin.flow.component.button.ButtonVariant;
-import com.vaadin.flow.component.html.H1;
-import com.vaadin.flow.component.orderedlayout.VerticalLayout;
-import com.vaadin.flow.router.Route;
-import org.axonframework.extensions.multitenancy.components.TenantDescriptor;
-import org.axonframework.extensions.multitenancy.components.TenantProvider;
-
-@Route
-public class MainView extends VerticalLayout {
-
- public MainView(TenantProvider tenantProvider) {
- VerticalLayout layout = new VerticalLayout();
-
- H1 title = new H1("Login");
- layout.add(title);
-
- layout.add(new Text("Available tenants:"));
-
- tenantProvider
- .getTenants()
- .stream()
- .map(TenantDescriptor::tenantId)
- .map(name -> new Button(name, event -> UI.getCurrent().navigate("action/" + name)))
- .forEach(layout::add);
-
- layout.setAlignItems(Alignment.CENTER);
-
-
- Button register = new Button("Or register new...", event -> UI.getCurrent().navigate("register"));
- register.addThemeVariants(ButtonVariant.LUMO_CONTRAST);
- layout.add(register);
-
- add(layout);
- }
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/ui/RegisterView.java b/multitenancy/src/main/java/io/axoniq/multitenancy/ui/RegisterView.java
deleted file mode 100644
index 57b6caec..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/ui/RegisterView.java
+++ /dev/null
@@ -1,69 +0,0 @@
-package io.axoniq.multitenancy.ui;
-
-import com.vaadin.flow.component.UI;
-import com.vaadin.flow.component.button.Button;
-import com.vaadin.flow.component.button.ButtonVariant;
-import com.vaadin.flow.component.html.H1;
-import com.vaadin.flow.component.notification.Notification;
-import com.vaadin.flow.component.notification.NotificationVariant;
-import com.vaadin.flow.component.orderedlayout.VerticalLayout;
-import com.vaadin.flow.component.textfield.TextField;
-import com.vaadin.flow.router.Route;
-import io.axoniq.multitenancy.web.CreateTenantService;
-
-@Route("/register")
-public class RegisterView extends VerticalLayout {
-
- public RegisterView(CreateTenantService createTenantService) {
- H1 title = new H1("Register tenant");
-
- VerticalLayout layout = new VerticalLayout();
- layout.setAlignItems(Alignment.CENTER);
- TextField tenantName = new TextField("Tenant name");
-
- layout.add(title);
-
- layout.add(tenantName);
-
- Button submitButton = new Button("Register");
- submitButton.addThemeVariants(ButtonVariant.LUMO_PRIMARY);
-
- layout.add(submitButton);
-
- submitButton.getStyle().set("padding", "25px");
-
- add(layout);
-
-
- submitButton.addClickListener(e -> {
- String tenantNameValue = tenantName.getValue();
-
- createTenantService
- .createTenant("tenant-" + tenantNameValue, "default", true)
- .exceptionally(err -> {
- showError(err.getMessage());
- return null;
- }).join();
-
- showSuccess(tenantNameValue);
-
- try {
- Thread.sleep(2500);
- } catch (InterruptedException ex) {
- ex.printStackTrace();
- }
-
- new Button("redirect", event -> UI.getCurrent().navigate("")).click();
- });
- }
-
- private void showSuccess(String tenantName) {
- Notification notification = Notification.show("Tenant created with username: tenant-" + tenantName);
- notification.addThemeVariants(NotificationVariant.LUMO_SUCCESS);
- }
-
- private void showError(String error) {
- Notification notification = Notification.show("Error has occurred: " + error);
- notification.addThemeVariants(NotificationVariant.LUMO_ERROR);
- }
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/ui/UIConfig.java b/multitenancy/src/main/java/io/axoniq/multitenancy/ui/UIConfig.java
deleted file mode 100644
index e458fddf..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/ui/UIConfig.java
+++ /dev/null
@@ -1,10 +0,0 @@
-package io.axoniq.multitenancy.ui;
-
-import com.vaadin.flow.component.page.AppShellConfigurator;
-import com.vaadin.flow.theme.Theme;
-import com.vaadin.flow.theme.lumo.Lumo;
-
-@Theme(variant = Lumo.DARK)
-class UIConfig implements AppShellConfigurator {
-
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/web/CreateTenantService.java b/multitenancy/src/main/java/io/axoniq/multitenancy/web/CreateTenantService.java
deleted file mode 100644
index 574f7d1f..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/web/CreateTenantService.java
+++ /dev/null
@@ -1,61 +0,0 @@
-/*
- * Copyright (c) 2020-2020. AxonIQ
- *
- * Licensed under the Apache License, Version 2.0 (the "License")
- * you may not use this file except in compliance with the License.
- * You may obtain a copy of the License at
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package io.axoniq.multitenancy.web;
-
-import io.axoniq.axonserver.connector.AxonServerConnection;
-import io.axoniq.axonserver.grpc.admin.CreateContextRequest;
-import org.axonframework.axonserver.connector.AxonServerConnectionManager;
-import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseBuilder;
-import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseType;
-import org.springframework.stereotype.Service;
-
-import java.util.concurrent.CompletableFuture;
-
-@Service
-public class CreateTenantService {
-
- private final AxonServerConnectionManager axonServerConnectionManager;
-
-
- public CreateTenantService(AxonServerConnectionManager axonServerConnectionManager) {
- this.axonServerConnectionManager = axonServerConnectionManager;
- }
-
- public CompletableFuture createTenant(String tenantName,
- String replicationGroup,
- boolean initializeSchema) {
- AxonServerConnection admin = axonServerConnectionManager.getConnection("_admin");
- CreateContextRequest createContextRequest = CreateContextRequest.newBuilder()
- .setName(tenantName)
- .setReplicationGroupName(replicationGroup)
- .build();
-
- return CompletableFuture.completedFuture(null)
- .thenRun(() -> {
- if (initializeSchema) {
- new EmbeddedDatabaseBuilder()
- .setType(EmbeddedDatabaseType.H2)
- .setName(tenantName + ";INIT=create " +
- "schema if not exists " +
- "schema_a\\;create schema if not exists schema_b;" +
- "DB_CLOSE_DELAY=-1;")
- .addScript("db/migration/V0.sql")
- .build();
- }
- })
- .thenCompose(s -> admin.adminChannel().createContext(createContextRequest));
- }
-}
diff --git a/multitenancy/src/main/java/io/axoniq/multitenancy/web/MessageService.java b/multitenancy/src/main/java/io/axoniq/multitenancy/web/MessageService.java
deleted file mode 100644
index 22d18126..00000000
--- a/multitenancy/src/main/java/io/axoniq/multitenancy/web/MessageService.java
+++ /dev/null
@@ -1,177 +0,0 @@
-package io.axoniq.multitenancy.web;
-
-import io.axoniq.multitenancy.api.AddFundsCommand;
-import io.axoniq.multitenancy.api.CardIssuedEvent;
-import io.axoniq.multitenancy.api.FindAllCardsQuery;
-import io.axoniq.multitenancy.api.FindCardQuery;
-import io.axoniq.multitenancy.api.GiftCardRecord;
-import io.axoniq.multitenancy.api.IssueCardCommand;
-import org.axonframework.commandhandling.CommandMessage;
-import org.axonframework.commandhandling.GenericCommandMessage;
-import org.axonframework.commandhandling.gateway.CommandGateway;
-import org.axonframework.eventhandling.EventBus;
-import org.axonframework.eventhandling.EventMessage;
-import org.axonframework.eventhandling.GenericEventMessage;
-import org.axonframework.eventhandling.scheduling.EventScheduler;
-import org.axonframework.messaging.MetaData;
-import org.axonframework.messaging.responsetypes.ResponseTypes;
-import org.axonframework.queryhandling.GenericQueryMessage;
-import org.axonframework.queryhandling.QueryGateway;
-import org.axonframework.queryhandling.QueryMessage;
-import org.axonframework.queryhandling.SubscriptionQueryResult;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import org.springframework.stereotype.Service;
-import reactor.core.publisher.Mono;
-
-import java.lang.invoke.MethodHandles;
-import java.time.Duration;
-import java.time.Instant;
-import java.util.Collections;
-import java.util.List;
-import java.util.Optional;
-import java.util.UUID;
-
-import static java.time.temporal.ChronoUnit.MINUTES;
-import static org.axonframework.commandhandling.GenericCommandMessage.asCommandMessage;
-import static org.axonframework.extensions.multitenancy.autoconfig.TenantConfiguration.TENANT_CORRELATION_KEY;
-
-@Service
-public class MessageService {
-
- private final static Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
-
- private final CommandGateway commandGateway;
- private final QueryGateway queryGateway;
- private final EventBus eventBus;
-
- private final EventScheduler eventScheduler;
-
- public MessageService(CommandGateway commandGateway, QueryGateway queryGateway, EventBus eventBus, EventScheduler eventScheduler) {
- this.commandGateway = commandGateway;
- this.queryGateway = queryGateway;
- this.eventBus = eventBus;
- this.eventScheduler = eventScheduler;
- }
-
- public void sendQueries(String tenantName) {
- for (int i = 0; i < 50; i++) {
- logger.info("Sending {}/{} queries delayed by {}", i + 1, 50, 50);
-
- QueryMessage> query =
- new GenericQueryMessage<>(
- new FindAllCardsQuery(),
- ResponseTypes.multipleInstancesOf(GiftCardRecord.class)
- ).withMetaData(Collections.singletonMap(TENANT_CORRELATION_KEY, tenantName));
-
- queryGateway.query(query, ResponseTypes.multipleInstancesOf(GiftCardRecord.class));
- }
- }
-
- public Mono subscriptionQuery(String tenantName) {
- final String giftCardId = UUID.randomUUID().toString();
- QueryMessage> query =
- new GenericQueryMessage<>(
- new FindCardQuery(giftCardId),
- ResponseTypes.optionalInstanceOf(GiftCardRecord.class)
- ).withMetaData(Collections.singletonMap(TENANT_CORRELATION_KEY, tenantName));
-
- SubscriptionQueryResult, GiftCardRecord> queryResult = queryGateway.subscriptionQuery(
- query,
- ResponseTypes.optionalInstanceOf(GiftCardRecord.class),
- ResponseTypes.instanceOf(GiftCardRecord.class));
-
- return sendAndReturnUpdate(new IssueCardCommand(giftCardId, 50), queryResult, tenantName)
- .map(GiftCardRecord::id);
- }
-
- public Mono sendAndReturnUpdate(Object command, SubscriptionQueryResult, U> result, String tenantName) {
- MetaData tenantMetaData = MetaData.with(TENANT_CORRELATION_KEY, tenantName);
- CommandMessageorg.testcontainers
- junit-jupiter
+ testcontainers-junit-jupitertest
diff --git a/snapshots/src/main/java/io/axoniq/dev/samples/SnapshottingApplication.java b/snapshots/src/main/java/io/axoniq/dev/samples/SnapshottingApplication.java
index 8e778a87..a20c21ef 100644
--- a/snapshots/src/main/java/io/axoniq/dev/samples/SnapshottingApplication.java
+++ b/snapshots/src/main/java/io/axoniq/dev/samples/SnapshottingApplication.java
@@ -1,14 +1,7 @@
package io.axoniq.dev.samples;
-import org.axonframework.eventsourcing.EventCountSnapshotTriggerDefinition;
-import org.axonframework.eventsourcing.SnapshotTriggerDefinition;
-import org.axonframework.eventsourcing.Snapshotter;
-import org.axonframework.serialization.Serializer;
-import org.axonframework.serialization.json.JacksonSerializer;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
-import org.springframework.context.annotation.Bean;
-import org.springframework.context.annotation.Primary;
@SpringBootApplication
public class SnapshottingApplication {
@@ -16,39 +9,4 @@ public class SnapshottingApplication {
public static void main(String[] args) {
SpringApplication.run(SnapshottingApplication.class, args);
}
-
- /**
- * This {@link SnapshotTriggerDefinition} is responsible to trigger the Snapshot. In this implementation, it creates
- * a Snapshot every 5 events, including the Snapshot itself. eg: when we hit the 5th event, a snapshot will be
- * created (5 events). When we hit the 9th event, another snapshot will be created (1 snapshot event + 4 events) and
- * so on.
- *
- * @param snapshotter The default {@link Snapshotter} provided by Axon.
- * @return the configured {@link SnapshotTriggerDefinition}.
- */
- @Bean
- public SnapshotTriggerDefinition mySnapshotTriggerDefinition(Snapshotter snapshotter) {
- return new EventCountSnapshotTriggerDefinition(snapshotter, 5);
- }
-
- /**
- * Construct the main {@link Serializer} used by this sample application to be a {@link JacksonSerializer}.
- *
- * By making this the {@link Primary} instance, Axon Framework will use this {@code Serializer} for commands,
- * events, queries, but more importantly, also snapshots. Since (by default) Axon Framework serializes the state of
- * the aggregate as the snapshot, this means the {@code io.axoniq.dev.samples.command.MyEntityAggregate} needs to be
- * serializable by Jackson.
- *
- * To support Jackson de-/serialization for this projects aggregate, we added the
- * {@link com.fasterxml.jackson.annotation.JsonGetter} and {@link com.fasterxml.jackson.annotation.JsonSetter}
- * annotation to getter and setter methods. Note that the getters/setters are made package-private on purpose, as
- * they should only be used by the serializer.
- *
- * @return The main {@link Serializer} used by this sample application to be a {@link JacksonSerializer}.
- */
- @Bean
- @Primary
- public Serializer serializer() {
- return JacksonSerializer.defaultSerializer();
- }
}
diff --git a/snapshots/src/main/java/io/axoniq/dev/samples/api/CreateMyEntityCommand.java b/snapshots/src/main/java/io/axoniq/dev/samples/api/CreateMyEntityCommand.java
index 64ed424a..a85d7318 100644
--- a/snapshots/src/main/java/io/axoniq/dev/samples/api/CreateMyEntityCommand.java
+++ b/snapshots/src/main/java/io/axoniq/dev/samples/api/CreateMyEntityCommand.java
@@ -1,9 +1,9 @@
package io.axoniq.dev.samples.api;
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
+import org.axonframework.modelling.annotation.TargetEntityId;
public record CreateMyEntityCommand(
- @TargetAggregateIdentifier String entityId,
+ @TargetEntityId String entityId,
String name
) {
diff --git a/snapshots/src/main/java/io/axoniq/dev/samples/api/RenameMyEntityCommand.java b/snapshots/src/main/java/io/axoniq/dev/samples/api/RenameMyEntityCommand.java
index 274e6a62..2e464ea9 100644
--- a/snapshots/src/main/java/io/axoniq/dev/samples/api/RenameMyEntityCommand.java
+++ b/snapshots/src/main/java/io/axoniq/dev/samples/api/RenameMyEntityCommand.java
@@ -1,9 +1,9 @@
package io.axoniq.dev.samples.api;
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
+import org.axonframework.modelling.annotation.TargetEntityId;
public record RenameMyEntityCommand(
- @TargetAggregateIdentifier String entityId,
+ @TargetEntityId String entityId,
String name
) {
diff --git a/snapshots/src/main/java/io/axoniq/dev/samples/command/MyEntityAggregate.java b/snapshots/src/main/java/io/axoniq/dev/samples/command/MyEntity.java
similarity index 65%
rename from snapshots/src/main/java/io/axoniq/dev/samples/command/MyEntityAggregate.java
rename to snapshots/src/main/java/io/axoniq/dev/samples/command/MyEntity.java
index 1acf15ee..2234bd6e 100644
--- a/snapshots/src/main/java/io/axoniq/dev/samples/command/MyEntityAggregate.java
+++ b/snapshots/src/main/java/io/axoniq/dev/samples/command/MyEntity.java
@@ -6,43 +6,43 @@
import io.axoniq.dev.samples.api.MyEntityCreatedEvent;
import io.axoniq.dev.samples.api.MyEntityRenamedEvent;
import io.axoniq.dev.samples.api.RenameMyEntityCommand;
-import org.axonframework.commandhandling.CommandHandler;
-import org.axonframework.eventsourcing.EventSourcingHandler;
-import org.axonframework.modelling.command.AggregateIdentifier;
-import org.axonframework.spring.stereotype.Aggregate;
+import org.axonframework.eventsourcing.annotation.EventSourcingHandler;
+import org.axonframework.eventsourcing.annotation.Snapshotting;
+import org.axonframework.eventsourcing.annotation.reflection.EntityCreator;
+import org.axonframework.extension.spring.stereotype.EventSourced;
+import org.axonframework.messaging.commandhandling.annotation.CommandHandler;
+import org.axonframework.messaging.eventhandling.gateway.EventAppender;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.lang.invoke.MethodHandles;
-import static org.axonframework.modelling.command.AggregateLifecycle.apply;
-
-@Aggregate(snapshotTriggerDefinition = "mySnapshotTriggerDefinition")
-class MyEntityAggregate {
+@EventSourced(tagKey = "MyEntity")
+@Snapshotting(afterEvents = 5)
+class MyEntity {
private static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
- @AggregateIdentifier
private String entityId;
private String name;
@CommandHandler
- public MyEntityAggregate(CreateMyEntityCommand command) {
+ public static void handle(CreateMyEntityCommand command, EventAppender appender) {
logger.info("[CreateMyEntityCommand] Entity with id [{}] and name [{}] created.",
command.entityId(), command.name());
- apply(new MyEntityCreatedEvent(command.entityId(), command.name()));
+ appender.append(new MyEntityCreatedEvent(command.entityId(), command.name()));
}
@CommandHandler
- public void on(RenameMyEntityCommand command) {
+ public void on(RenameMyEntityCommand command, EventAppender appender) {
logger.info("[RenameMyEntityCommand] Entity with id [{}] and name [{}] updated.",
command.entityId(), command.name());
if (name.equals(command.name())) {
throw new IllegalArgumentException("New name can not be the same as current name.");
}
- apply(new MyEntityRenamedEvent(command.entityId(), command.name()));
+ appender.append(new MyEntityRenamedEvent(command.entityId(), command.name()));
}
@EventSourcingHandler
@@ -60,7 +60,8 @@ public void on(MyEntityRenamedEvent event) {
logger.info("[MyEntityRenamedEvent] Entity with id [{}] being event sourced.", event.entityId());
}
- // Since the main Serializer is a JacksonSerializer, the constructed Snapshot will also be serialized through Jackson
+ // The general Converter defaults to a JacksonConverter, so the constructed Snapshot will also be
+ // (de)serialized through Jackson.
@JsonGetter
String getEntityId() {
return entityId;
@@ -81,7 +82,8 @@ void setName(String name) {
this.name = name;
}
- public MyEntityAggregate() {
+ @EntityCreator
+ public MyEntity() {
// Required by Axon Framework
}
}
diff --git a/snapshots/src/main/java/io/axoniq/dev/samples/rest/CommandController.java b/snapshots/src/main/java/io/axoniq/dev/samples/rest/CommandController.java
index 4c315c73..936ba8dd 100644
--- a/snapshots/src/main/java/io/axoniq/dev/samples/rest/CommandController.java
+++ b/snapshots/src/main/java/io/axoniq/dev/samples/rest/CommandController.java
@@ -2,7 +2,7 @@
import io.axoniq.dev.samples.api.CreateMyEntityCommand;
import io.axoniq.dev.samples.api.RenameMyEntityCommand;
-import org.axonframework.commandhandling.gateway.CommandGateway;
+import org.axonframework.messaging.commandhandling.gateway.CommandGateway;
import org.springframework.web.bind.annotation.PatchMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
@@ -25,13 +25,13 @@ public CommandController(CommandGateway commandGateway) {
@PostMapping("/{id}")
public CompletableFuture createMyEntity(@PathVariable("id") String entityId,
@RequestParam("name") String name) {
- return commandGateway.send(new CreateMyEntityCommand(entityId, name));
+ return commandGateway.send(new CreateMyEntityCommand(entityId, name), Void.class);
}
@PatchMapping("/{id}")
public CompletableFuture renameMyEntity(@PathVariable("id") String entityId,
@RequestParam("name") String name) {
- return commandGateway.send(new RenameMyEntityCommand(entityId, name));
+ return commandGateway.send(new RenameMyEntityCommand(entityId, name), Void.class);
}
}
diff --git a/stateful-event-handler/pom.xml b/stateful-event-handler/pom.xml
index 4175c3f6..0b3081bc 100644
--- a/stateful-event-handler/pom.xml
+++ b/stateful-event-handler/pom.xml
@@ -21,7 +21,7 @@
axon-messaging
- org.axonframework
+ org.axonframework.extensions.springaxon-spring-boot-starter
diff --git a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/config/OrderProcessorConfig.java b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/config/OrderProcessorConfig.java
index 26276954..00f78022 100644
--- a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/config/OrderProcessorConfig.java
+++ b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/config/OrderProcessorConfig.java
@@ -1,8 +1,6 @@
package io.axoniq.dev.samples.config;
-import org.axonframework.config.ConfigurerModule;
-import org.axonframework.eventhandling.TrackingEventProcessorConfiguration;
-import org.axonframework.messaging.StreamableMessageSource;
+import org.axonframework.extension.spring.config.EventProcessorDefinition;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@@ -10,11 +8,19 @@
public class OrderProcessorConfig {
@Bean
- public ConfigurerModule configureOrderProcessor() {
- TrackingEventProcessorConfiguration tepConfig =
- TrackingEventProcessorConfiguration.forSingleThreadedProcessing()
- .andInitialTrackingToken(StreamableMessageSource::createHeadToken);
- return configurer -> configurer.eventProcessing()
- .registerTrackingEventProcessorConfiguration("OrderProcessor", c -> tepConfig);
+ public EventProcessorDefinition orderProcessorDefinition() {
+ // A PooledStreamingEventProcessor is inherently multi-segment/multi-threaded (16 segments by
+ // default). Restricting it to a single initial segment means there is only ever one segment
+ // to claim, so only one worker thread will ever be processing for this processor - reproducing
+ // AF4's TrackingEventProcessorConfiguration#forSingleThreadedProcessing() semantics.
+ //
+ // AF4's StreamableMessageSource#createHeadToken() maps directly onto AF5's
+ // TrackingTokenSource#latestToken(...): both create a token pointing at the current end of the
+ // stream, so only events published after start-up are processed and earlier events (which could
+ // otherwise trigger unwanted side effects, such as re-sending commands) are skipped.
+ return EventProcessorDefinition.pooledStreamingMatching("OrderProcessor")
+ .customized(config -> config
+ .initialSegmentCount(1)
+ .initialToken(source -> source.latestToken(null)));
}
}
diff --git a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/order/api/CompleteOrderProcessCommand.java b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/order/api/CompleteOrderProcessCommand.java
index 64c307dd..ddcbfdb4 100644
--- a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/order/api/CompleteOrderProcessCommand.java
+++ b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/order/api/CompleteOrderProcessCommand.java
@@ -1,10 +1,10 @@
package io.axoniq.dev.samples.order.api;
import io.axoniq.dev.samples.uuid.OrderId;
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
+import org.axonframework.modelling.annotation.TargetEntityId;
public record CompleteOrderProcessCommand(
- @TargetAggregateIdentifier OrderId orderId,
+ @TargetEntityId OrderId orderId,
boolean isPaid,
boolean orderIsDelivered
) {
diff --git a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/orderprocessor/OrderProcessor.java b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/orderprocessor/OrderProcessor.java
index 9599856b..8535c850 100644
--- a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/orderprocessor/OrderProcessor.java
+++ b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/orderprocessor/OrderProcessor.java
@@ -12,14 +12,14 @@
import io.axoniq.dev.samples.uuid.PaymentId;
import io.axoniq.dev.samples.uuid.ShipmentId;
import io.axoniq.dev.samples.uuid.UUIDProvider;
-import org.axonframework.commandhandling.gateway.CommandGateway;
-import org.axonframework.config.ProcessingGroup;
-import org.axonframework.eventhandling.EventHandler;
+import org.axonframework.messaging.commandhandling.gateway.CommandGateway;
+import org.axonframework.messaging.core.annotation.Namespace;
+import org.axonframework.messaging.eventhandling.annotation.EventHandler;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
@Component
-@ProcessingGroup("OrderProcessor")
+@Namespace("OrderProcessor")
class OrderProcessor {
private final CommandGateway commandGateway;
diff --git a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/payment/api/PayOrderCommand.java b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/payment/api/PayOrderCommand.java
index f66d2c4e..60ae146a 100644
--- a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/payment/api/PayOrderCommand.java
+++ b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/payment/api/PayOrderCommand.java
@@ -1,10 +1,10 @@
package io.axoniq.dev.samples.payment.api;
import io.axoniq.dev.samples.uuid.PaymentId;
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
+import org.axonframework.modelling.annotation.TargetEntityId;
public record PayOrderCommand(
- @TargetAggregateIdentifier PaymentId paymentId
+ @TargetEntityId PaymentId paymentId
) {
}
diff --git a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/shipment/api/CancelShipmentCommand.java b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/shipment/api/CancelShipmentCommand.java
index fc582578..c937ce72 100644
--- a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/shipment/api/CancelShipmentCommand.java
+++ b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/shipment/api/CancelShipmentCommand.java
@@ -1,10 +1,10 @@
package io.axoniq.dev.samples.shipment.api;
import io.axoniq.dev.samples.uuid.ShipmentId;
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
+import org.axonframework.modelling.annotation.TargetEntityId;
public record CancelShipmentCommand(
- @TargetAggregateIdentifier ShipmentId shipmentId
+ @TargetEntityId ShipmentId shipmentId
) {
}
diff --git a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/shipment/api/ShipOrderCommand.java b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/shipment/api/ShipOrderCommand.java
index 50bf62f1..f2ec4a8a 100644
--- a/stateful-event-handler/src/main/java/io/axoniq/dev/samples/shipment/api/ShipOrderCommand.java
+++ b/stateful-event-handler/src/main/java/io/axoniq/dev/samples/shipment/api/ShipOrderCommand.java
@@ -1,10 +1,10 @@
package io.axoniq.dev.samples.shipment.api;
import io.axoniq.dev.samples.uuid.ShipmentId;
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
+import org.axonframework.modelling.annotation.TargetEntityId;
public record ShipOrderCommand(
- @TargetAggregateIdentifier ShipmentId shipmentId
+ @TargetEntityId ShipmentId shipmentId
) {
}
diff --git a/subscription-query-rest/README.md b/subscription-query-rest/README.md
index dfa9d14f..79ae139d 100644
--- a/subscription-query-rest/README.md
+++ b/subscription-query-rest/README.md
@@ -6,12 +6,13 @@ projection right away, instead of listening for updates on a different endpoint.
There are two issues that needs to be address in this case:
1. We need to subscribe for updates before we send a command, that’s the only way to be sure we will not miss any
- updates. Sending commands first and then subscribing for updates will result in race conditions! The simple trick is
- to subscribe for the initial result first (even if we don’t need it). Let’s call it virtual initial result. This will
- open the Subscription query, which will buffer all updates that arrive at this point on. Since we now have a buffer
- for updates, we can send a command and after the command has sent we can subscribe to updates flux. If an update
- arrives after sending a command and before we are subscribed for updates, we will read it automatically from the
- buffer, therefore we are sure we will not miss any updates.
+ updates. Sending commands first and then subscribing for updates will result in race conditions! Axon Framework's
+ `QueryGateway#subscriptionQuery(...)` returns a single reactive-streams `Publisher` that combines the (here virtual,
+ empty) initial result and every subsequent update: the query is only sent, and the update buffer only opened, once
+ that `Publisher` is subscribed to. The simple trick is therefore to subscribe to it first, and only dispatch the
+ command once that subscription is established (e.g. from a `doOnSubscribe` callback). If an update arrives right
+ after sending the command and before we start consuming from the `Publisher`, we will still read it from the
+ buffer that was already open, therefore we are sure we will not miss any updates.
2. We need to read our own writes, multiple updates/events could be dispatch at the same time, we can’t guarantee order
and which one will arrive first. Without some kind of correlation, we will easily get into trouble and get someone
diff --git a/subscription-query-rest/pom.xml b/subscription-query-rest/pom.xml
index 65528137..8dc7be4e 100644
--- a/subscription-query-rest/pom.xml
+++ b/subscription-query-rest/pom.xml
@@ -19,7 +19,7 @@
- org.axonframework
+ org.axonframework.extensions.springaxon-spring-boot-starter
@@ -49,7 +49,7 @@
org.testcontainers
- junit-jupiter
+ testcontainers-junit-jupiter
\ No newline at end of file
diff --git a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/CommandController.java b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/CommandController.java
index e4afbf6f..c9ac804e 100644
--- a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/CommandController.java
+++ b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/CommandController.java
@@ -3,14 +3,16 @@
import io.axoniq.dev.samples.api.CreateMyEntityCommand;
import io.axoniq.dev.samples.api.GetMyEntityByCorrelationIdQuery;
import io.axoniq.dev.samples.query.MyEntity;
-import org.axonframework.commandhandling.CommandMessage;
-import org.axonframework.commandhandling.GenericCommandMessage;
-import org.axonframework.commandhandling.gateway.CommandGateway;
-import org.axonframework.queryhandling.QueryGateway;
-import org.axonframework.queryhandling.SubscriptionQueryResult;
+import org.axonframework.messaging.commandhandling.CommandMessage;
+import org.axonframework.messaging.commandhandling.GenericCommandMessage;
+import org.axonframework.messaging.commandhandling.gateway.CommandGateway;
+import org.axonframework.messaging.core.MessageType;
+import org.axonframework.messaging.queryhandling.gateway.QueryGateway;
+import org.reactivestreams.Publisher;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RestController;
+import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.time.Duration;
@@ -29,31 +31,30 @@ public CommandController(CommandGateway commandGateway, QueryGateway queryGatewa
@PostMapping("/entities/{id}")
public Mono myApi(@PathVariable("id") String entityId) {
- /* We are wrapping command into GenericCommandMessage, so we can get its identifier (correlation id) */
- CommandMessage command = GenericCommandMessage.asCommandMessage(new CreateMyEntityCommand(entityId));
+ /* We are wrapping the command into a GenericCommandMessage, so we can get its identifier (correlation id) */
+ CommandMessage command = new GenericCommandMessage(new MessageType(CreateMyEntityCommand.class),
+ new CreateMyEntityCommand(entityId));
- /* With command identifier we can now subscribe for updates that this command produced */
- GetMyEntityByCorrelationIdQuery query = new GetMyEntityByCorrelationIdQuery(command.getIdentifier());
+ /* With the command identifier we can now subscribe for updates that this command produced */
+ GetMyEntityByCorrelationIdQuery query = new GetMyEntityByCorrelationIdQuery(command.identifier());
- /* since we don't care about initial result, we mark it as Void.class */
- SubscriptionQueryResult response = queryGateway.subscriptionQuery(query,
- Void.class,
- MyEntity.class);
- return sendAndReturnUpdate(command, response)
+ /* Axon Framework 5 merges the "virtual" initial result and the updates into a single Publisher,
+ so there is only one responseType now instead of separate initial/update types */
+ Publisher updates = queryGateway.subscriptionQuery(query, MyEntity.class);
+
+ return sendAndReturnUpdate(command, updates)
.map(MyEntity::id);
}
- public Mono sendAndReturnUpdate(Object command, SubscriptionQueryResult, U> result) {
- /* The trick here is to subscribe to initial results first, even it does not return any result
- Subscribing to initialResult creates a buffer for updates, even that we didn't subscribe for updates yet
- they will wait for us in buffer, after this we can safely send command, and then subscribe to updates */
- return Mono.when(result.initialResult())
- .then(Mono.fromCompletionStage(() -> commandGateway.send(command)))
- .thenMany(result.updates())
+ public Mono sendAndReturnUpdate(Object command, Publisher updates) {
+ /* The trick here is to subscribe to the subscription query's Publisher first: subscribing is what sends
+ the query and opens the buffer for updates. We hook the command dispatch into doOnSubscribe so it only
+ fires once that buffer is open, guaranteeing we cannot miss the update it produces. Cancelling (closing)
+ the subscription query is handled automatically once next() receives its element or the timeout fires. */
+ return Flux.from(updates)
+ .doOnSubscribe(subscription -> commandGateway.send(command))
.timeout(Duration.ofSeconds(5))
- .next()
- .doFinally(unused -> result.cancel());
- /* dont forget to close subscription query on the end and add a timeout */
+ .next();
}
}
diff --git a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/CreateMyEntityCommand.java b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/CreateMyEntityCommand.java
index d891bb9b..356998d4 100644
--- a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/CreateMyEntityCommand.java
+++ b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/CreateMyEntityCommand.java
@@ -1,9 +1,9 @@
package io.axoniq.dev.samples.api;
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
+import org.axonframework.modelling.annotation.TargetEntityId;
public record CreateMyEntityCommand(
- @TargetAggregateIdentifier String entityId
+ @TargetEntityId String entityId
) {
}
diff --git a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/MyEntityCreatedEvent.java b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/MyEntityCreatedEvent.java
index be9b1090..72625a31 100644
--- a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/MyEntityCreatedEvent.java
+++ b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/MyEntityCreatedEvent.java
@@ -1,7 +1,9 @@
package io.axoniq.dev.samples.api;
+import org.axonframework.eventsourcing.annotation.EventTag;
+
public record MyEntityCreatedEvent(
- String entityId
+ @EventTag(key = "MyEntity") String entityId
) {
}
diff --git a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/ValidateMyEntityCommand.java b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/ValidateMyEntityCommand.java
index 9ec1f288..5b927fd9 100644
--- a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/ValidateMyEntityCommand.java
+++ b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/api/ValidateMyEntityCommand.java
@@ -1,9 +1,9 @@
package io.axoniq.dev.samples.api;
-import org.axonframework.modelling.command.TargetAggregateIdentifier;
+import org.axonframework.modelling.annotation.TargetEntityId;
public record ValidateMyEntityCommand(
- @TargetAggregateIdentifier String entityId, String email
+ @TargetEntityId String entityId, String email
) {
}
diff --git a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/command/MyEntity.java b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/command/MyEntity.java
new file mode 100644
index 00000000..3a44a8b8
--- /dev/null
+++ b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/command/MyEntity.java
@@ -0,0 +1,30 @@
+package io.axoniq.dev.samples.command;
+
+import io.axoniq.dev.samples.api.CreateMyEntityCommand;
+import io.axoniq.dev.samples.api.MyEntityCreatedEvent;
+import org.axonframework.eventsourcing.annotation.EventSourcingHandler;
+import org.axonframework.eventsourcing.annotation.reflection.EntityCreator;
+import org.axonframework.extension.spring.stereotype.EventSourced;
+import org.axonframework.messaging.commandhandling.annotation.CommandHandler;
+import org.axonframework.messaging.eventhandling.gateway.EventAppender;
+
+@EventSourced(tagKey = "MyEntity")
+class MyEntity {
+
+ private String entityId;
+
+ @EntityCreator
+ public MyEntity() {
+ // Required by Axon Framework
+ }
+
+ @CommandHandler
+ public static void handle(CreateMyEntityCommand command, EventAppender eventAppender) {
+ eventAppender.append(new MyEntityCreatedEvent(command.entityId()));
+ }
+
+ @EventSourcingHandler
+ public void on(MyEntityCreatedEvent event) {
+ entityId = event.entityId();
+ }
+}
diff --git a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/command/MyEntityAggregate.java b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/command/MyEntityAggregate.java
deleted file mode 100644
index 9d76a264..00000000
--- a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/command/MyEntityAggregate.java
+++ /dev/null
@@ -1,31 +0,0 @@
-package io.axoniq.dev.samples.command;
-
-import io.axoniq.dev.samples.api.CreateMyEntityCommand;
-import io.axoniq.dev.samples.api.MyEntityCreatedEvent;
-import org.axonframework.commandhandling.CommandHandler;
-import org.axonframework.eventsourcing.EventSourcingHandler;
-import org.axonframework.modelling.command.AggregateIdentifier;
-import org.axonframework.spring.stereotype.Aggregate;
-
-import static org.axonframework.modelling.command.AggregateLifecycle.apply;
-
-@Aggregate
-class MyEntityAggregate {
-
- @AggregateIdentifier
- private String entityId;
-
- public MyEntityAggregate() {
- // Required by Axon Framework
- }
-
- @CommandHandler
- public MyEntityAggregate(CreateMyEntityCommand command) {
- apply(new MyEntityCreatedEvent(command.entityId()));
- }
-
- @EventSourcingHandler
- public void on(MyEntityCreatedEvent event) {
- entityId = event.entityId();
- }
-}
diff --git a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/query/MyEntityProjection.java b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/query/MyEntityProjection.java
index 02faad27..766c4f46 100644
--- a/subscription-query-rest/src/main/java/io/axoniq/dev/samples/query/MyEntityProjection.java
+++ b/subscription-query-rest/src/main/java/io/axoniq/dev/samples/query/MyEntityProjection.java
@@ -2,10 +2,10 @@
import io.axoniq.dev.samples.api.GetMyEntityByCorrelationIdQuery;
import io.axoniq.dev.samples.api.MyEntityCreatedEvent;
-import org.axonframework.eventhandling.EventHandler;
-import org.axonframework.messaging.annotation.MetaDataValue;
-import org.axonframework.queryhandling.QueryHandler;
-import org.axonframework.queryhandling.QueryUpdateEmitter;
+import org.axonframework.messaging.core.annotation.MetadataValue;
+import org.axonframework.messaging.eventhandling.annotation.EventHandler;
+import org.axonframework.messaging.queryhandling.QueryUpdateEmitter;
+import org.axonframework.messaging.queryhandling.annotation.QueryHandler;
import org.springframework.stereotype.Component;
import java.util.Optional;
@@ -13,12 +13,6 @@
@Component
class MyEntityProjection {
- private final QueryUpdateEmitter emitter;
-
- public MyEntityProjection(QueryUpdateEmitter emitter) {
- this.emitter = emitter;
- }
-
@QueryHandler
/* We are creating virtual initial result, doesn't need to return anything, but also do not return null */
public Optional on(GetMyEntityByCorrelationIdQuery query) {
@@ -26,7 +20,8 @@ public Optional on(GetMyEntityByCorrelationIdQuery query) {
}
@EventHandler
- public void on(MyEntityCreatedEvent event, @MetaDataValue("correlationId") String correlationId) {
+ public void on(MyEntityCreatedEvent event, @MetadataValue("correlationId") String correlationId,
+ QueryUpdateEmitter emitter) {
MyEntity entity = new MyEntity(event.entityId());
/* save your entity in your repository here */
diff --git a/subscription-query-rest/src/main/resources/application.properties b/subscription-query-rest/src/main/resources/application.properties
index 7ed4f46c..04a7a8ef 100644
--- a/subscription-query-rest/src/main/resources/application.properties
+++ b/subscription-query-rest/src/main/resources/application.properties
@@ -1 +1 @@
-axon.serializer.general=jackson
\ No newline at end of file
+axon.converter.general=jackson
\ No newline at end of file
diff --git a/subscription-query-streaming/README.md b/subscription-query-streaming/README.md
index 8ac29c88..97353ca9 100644
--- a/subscription-query-streaming/README.md
+++ b/subscription-query-streaming/README.md
@@ -19,9 +19,10 @@ First of these is the `EventPublisher` in the `commandmodel` package. The `Event
a `StreamUpdatedEvent` to spoof an active applications. It fills the `StreamUpdatedEvent` with `UUIDs`.
Secondly, the `ModelProjector` in the `querymodel` package is in charge of handling the `StreamUpdatedEvent`. It adds th
-e contents to a `List` of strings, and emits an update through the `QueryUpdateEmitter`. Next to that, a `@QueryHandler`
-annotated method is present for the `ModelQuery`. This query handler returns the entire list of updates. The combination
-of this query handler, and the update emission on the event handler provide an entry point for a subscription query.
+e contents to a `List` of strings, and emits an update through the `QueryUpdateEmitter`, which is injected as a
+parameter of the `@EventHandler` method. Next to that, a `@QueryHandler` annotated method is present for the
+`ModelQuery`. This query handler returns the entire list of updates. The combination of this query handler, and the
+update emission on the event handler provide an entry point for a subscription query.
Thirdly, the `QueryController` in the `ui` package provides an endpoint on `/app/updates`. This endpoint returns
a `Flux` of `ServerSentEvents`. The`ServerSentEvents` are filled with the result of a subscription query on
diff --git a/subscription-query-streaming/pom.xml b/subscription-query-streaming/pom.xml
index a5aa94b3..6527b7cc 100644
--- a/subscription-query-streaming/pom.xml
+++ b/subscription-query-streaming/pom.xml
@@ -17,7 +17,7 @@
- org.axonframework
+ org.axonframework.extensions.springaxon-spring-boot-starter
diff --git a/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/AxonConfig.java b/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/AxonConfig.java
index e9af7066..2a4403ee 100644
--- a/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/AxonConfig.java
+++ b/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/AxonConfig.java
@@ -1,10 +1,10 @@
package io.axoniq.dev.samples;
-import org.axonframework.commandhandling.gateway.CommandGateway;
-import org.axonframework.eventhandling.gateway.EventGateway;
-import org.axonframework.messaging.Message;
-import org.axonframework.messaging.interceptors.LoggingInterceptor;
-import org.axonframework.queryhandling.QueryGateway;
+import org.axonframework.messaging.commandhandling.gateway.CommandGateway;
+import org.axonframework.messaging.core.Message;
+import org.axonframework.messaging.core.interception.LoggingInterceptor;
+import org.axonframework.messaging.eventhandling.gateway.EventGateway;
+import org.axonframework.messaging.queryhandling.gateway.QueryGateway;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@@ -18,7 +18,7 @@
public class AxonConfig {
@Bean
- public LoggingInterceptor> loggingInterceptor() {
+ public LoggingInterceptor loggingInterceptor() {
return new LoggingInterceptor<>();
}
diff --git a/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/commandmodel/EventPublisher.java b/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/commandmodel/EventPublisher.java
index 5d626ef8..4296a5ff 100644
--- a/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/commandmodel/EventPublisher.java
+++ b/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/commandmodel/EventPublisher.java
@@ -1,7 +1,7 @@
package io.axoniq.dev.samples.commandmodel;
import io.axoniq.dev.samples.api.StreamUpdatedEvent;
-import org.axonframework.eventhandling.gateway.EventGateway;
+import org.axonframework.messaging.eventhandling.gateway.EventGateway;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
@@ -24,6 +24,8 @@ public EventPublisher(EventGateway eventGateway) {
@Scheduled(initialDelay = 1_000, fixedDelay = 6_000)
public void publishEvent() {
- eventGateway.publish(new StreamUpdatedEvent(UUID.randomUUID().toString()));
+ // AF5's EventGateway no longer exposes a plain publish(Object...) — the varargs overload now
+ // requires a ProcessingContext, which we don't have from a @Scheduled method, hence null.
+ eventGateway.publish(null, new StreamUpdatedEvent(UUID.randomUUID().toString()));
}
}
diff --git a/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/querymodel/ModelProjector.java b/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/querymodel/ModelProjector.java
index bb822aff..424d57d2 100644
--- a/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/querymodel/ModelProjector.java
+++ b/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/querymodel/ModelProjector.java
@@ -2,10 +2,10 @@
import io.axoniq.dev.samples.api.ModelQuery;
import io.axoniq.dev.samples.api.StreamUpdatedEvent;
-import org.axonframework.config.ProcessingGroup;
-import org.axonframework.eventhandling.EventHandler;
-import org.axonframework.queryhandling.QueryHandler;
-import org.axonframework.queryhandling.QueryUpdateEmitter;
+import org.axonframework.messaging.core.annotation.Namespace;
+import org.axonframework.messaging.eventhandling.annotation.EventHandler;
+import org.axonframework.messaging.queryhandling.QueryUpdateEmitter;
+import org.axonframework.messaging.queryhandling.annotation.QueryHandler;
import org.springframework.stereotype.Component;
import java.util.ArrayList;
@@ -18,19 +18,17 @@
* @author Steven van Beelen
*/
@Component
-@ProcessingGroup("model-projector")
+@Namespace("model-projector")
public class ModelProjector {
- private final QueryUpdateEmitter updateEmitter;
private final List updates;
- public ModelProjector(QueryUpdateEmitter updateEmitter) {
- this.updateEmitter = updateEmitter;
+ public ModelProjector() {
this.updates = new ArrayList<>();
}
@EventHandler
- public void on(StreamUpdatedEvent event) {
+ public void on(StreamUpdatedEvent event, QueryUpdateEmitter updateEmitter) {
updates.add(event.update());
updateEmitter.emit(ModelQuery.class, query -> true, event.update());
}
diff --git a/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/ui/QueryController.java b/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/ui/QueryController.java
index ecac856d..807cb727 100644
--- a/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/ui/QueryController.java
+++ b/subscription-query-streaming/src/main/java/io/axoniq/dev/samples/ui/QueryController.java
@@ -1,9 +1,9 @@
package io.axoniq.dev.samples.ui;
import io.axoniq.dev.samples.api.ModelQuery;
-import org.axonframework.messaging.responsetypes.ResponseTypes;
-import org.axonframework.queryhandling.QueryGateway;
-import org.axonframework.queryhandling.SubscriptionQueryResult;
+import org.axonframework.messaging.queryhandling.QueryResponseMessage;
+import org.axonframework.messaging.queryhandling.SubscriptionQueryUpdateMessage;
+import org.axonframework.messaging.queryhandling.gateway.QueryGateway;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.http.MediaType;
@@ -39,21 +39,23 @@ public QueryController(QueryGateway queryGateway) {
@CrossOrigin(exposedHeaders = "Access-Control-Allow-Origin")
@GetMapping(path = "/updates", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux> updates() {
- //noinspection resource
- SubscriptionQueryResult, String> result =
- queryGateway.subscriptionQuery(new ModelQuery(),
- ResponseTypes.multipleInstancesOf(String.class),
- ResponseTypes.instanceOf(String.class));
-
- Flux> sseStream = result.initialResult()
- .flatMapMany(Flux::fromIterable)
- .concatWith(result.updates())
- .doOnError(throwable -> logger.warn("something failed"))
- .map(update -> ServerSentEvent.builder()
- .event("update")
- .data(update)
- .build())
- .doFinally(signal -> result.close());
+ // Axon Framework 5's QueryGateway#subscriptionQuery combines the initial result and the updates into a
+ // single Publisher of one responseType, and there is no more SubscriptionQueryResult to close explicitly:
+ // the underlying subscription is closed automatically once the returned Flux is cancelled/disposed.
+ // The ModelQuery's initial result is a List, while every emitted update is a single String. To
+ // keep that shape, we fall back to the mapper-based subscriptionQuery overload, which lets us
+ // distinguish the initial result from an update via the message type, and flatten the initial
+ // List into individual elements ourselves (mirroring the old initialResult().flatMapMany(...)).
+ Flux> sseStream =
+ Flux.from(queryGateway.subscriptionQuery(new ModelQuery(), Object.class, QueryController::mapResponse))
+ .flatMap(response -> response instanceof List> initialResult
+ ? Flux.fromIterable(initialResult).cast(String.class)
+ : Flux.just((String) response))
+ .doOnError(throwable -> logger.warn("something failed"))
+ .map(update -> ServerSentEvent.builder()
+ .event("update")
+ .data(update)
+ .build());
// For Server Sent Events, the server doesn't get a close signal when the client closes the connection.
// Hence, we are left with a hanging stream in that case.
@@ -65,4 +67,10 @@ public Flux> updates() {
.build());
return Flux.merge(sseStream, heartbeatStream);
}
+
+ private static Object mapResponse(QueryResponseMessage response) {
+ return response instanceof SubscriptionQueryUpdateMessage
+ ? response.payloadAs(String.class)
+ : response.payloadAs(List.class);
+ }
}
diff --git a/subscription-query-streaming/src/main/resources/application.properties b/subscription-query-streaming/src/main/resources/application.properties
index 7f69874d..ec13b4aa 100644
--- a/subscription-query-streaming/src/main/resources/application.properties
+++ b/subscription-query-streaming/src/main/resources/application.properties
@@ -1,2 +1,2 @@
spring.application.name=Subscription Query Streaming
-axon.serializer.general=jackson
\ No newline at end of file
+axon.converter.general=jackson
\ No newline at end of file
diff --git a/upcaster/README.md b/upcaster/README.md
deleted file mode 100644
index 7606c831..00000000
--- a/upcaster/README.md
+++ /dev/null
@@ -1,40 +0,0 @@
-# Upcasters
-
-This module's intention is to give insights on how to create an upcaster and how to test it. This example uses a Json
-serializer.
-
-The FlightDelayedEvent is used as example. The original (null) revision looks like this:
-
-```json
-{
- "arrivalTime": "2021-05-27T15:06:10.629267",
- "flightId": "KL123",
- "origin": "LAX",
- "destination": "LON"
-}
-
-```
-
-Then a new requirement popped up to put the origin and destination into a separate object called `leg`. Please note,
-that this is just an example, and we do recommend keeping your event structure as flat as possible. The upcasted event
-should look like this:
-
-```json
-{
- "arrivalTime": "2021-05-27T15:06:10.629267",
- "flightId": "KL123",
- "leg": {
- "origin": "LAX",
- "destination": "LON"
- }
-}
-```
-
-You can find the implementation of the upcaster in
-the [FlightDelayedEventUpcaster](src/main/java/io/axoniq/dev/samples/upcaster/json/FlightDelayedEvent0_to_1Upcaster.java).
-The implementation of the test can be
-found [here](src/test/java/io/axoniq/dev/samples/upcaster/json/FlightDelayedEvent0_To_1UpcasterTest.java)
-
-To get this upcaster invoked on the event handler it should be added to
-the [EventUpcasterChainFactory](src/main/java/io/axoniq/dev/samples/upcaster/json/EventUpcasterChainFactory.java) or
-annotate it as a Spring component together with an Order annotation.
diff --git a/upcaster/pom.xml b/upcaster/pom.xml
deleted file mode 100644
index 28900fcd..00000000
--- a/upcaster/pom.xml
+++ /dev/null
@@ -1,44 +0,0 @@
-
-
-
- code-samples
- io.axoniq
- 0.0.2-SNAPSHOT
-
- 4.0.0
-
- upcaster
-
- Upcaster
- Module showing several Upcaster implementations
-
-
- 1.5.3
-
-
-
-
-
- org.axonframework
- axon-spring-boot-starter
-
-
-
- org.springframework.boot
- spring-boot-starter-web
-
-
-
- org.skyscreamer
- jsonassert
- ${jsonassert.version}
- test
-
-
- org.junit.jupiter
- junit-jupiter
-
-
-
\ No newline at end of file
diff --git a/upcaster/src/main/java/io/axoniq/dev/samples/api/AirportCode.java b/upcaster/src/main/java/io/axoniq/dev/samples/api/AirportCode.java
deleted file mode 100644
index 8b45460f..00000000
--- a/upcaster/src/main/java/io/axoniq/dev/samples/api/AirportCode.java
+++ /dev/null
@@ -1,12 +0,0 @@
-package io.axoniq.dev.samples.api;
-
-public enum AirportCode {
- AMS("Amsterdam"),
- LON("London"),
- PAR("Paris"),
- NYC("New York City"),
- LAX("Los Angeles");
-
- AirportCode(@SuppressWarnings("unused") String airportName) {
- }
-}
diff --git a/upcaster/src/main/java/io/axoniq/dev/samples/api/FlightDelayedEvent.java b/upcaster/src/main/java/io/axoniq/dev/samples/api/FlightDelayedEvent.java
deleted file mode 100644
index 845eaaf0..00000000
--- a/upcaster/src/main/java/io/axoniq/dev/samples/api/FlightDelayedEvent.java
+++ /dev/null
@@ -1,15 +0,0 @@
-package io.axoniq.dev.samples.api;
-
-import org.axonframework.serialization.Revision;
-
-import java.time.LocalDateTime;
-import java.util.Objects;
-
-@Revision("1.0")
-public record FlightDelayedEvent(
- String flightId,
- Leg leg,
- LocalDateTime arrivalTime
-) {
-
-}
diff --git a/upcaster/src/main/java/io/axoniq/dev/samples/api/Leg.java b/upcaster/src/main/java/io/axoniq/dev/samples/api/Leg.java
deleted file mode 100644
index abf1d1b4..00000000
--- a/upcaster/src/main/java/io/axoniq/dev/samples/api/Leg.java
+++ /dev/null
@@ -1,8 +0,0 @@
-package io.axoniq.dev.samples.api;
-
-public record Leg(
- AirportCode origin,
- AirportCode destination
-) {
-
-}
diff --git a/upcaster/src/main/java/io/axoniq/dev/samples/api/PassengerSeatAdjustedEvent.java b/upcaster/src/main/java/io/axoniq/dev/samples/api/PassengerSeatAdjustedEvent.java
deleted file mode 100644
index 169a7c09..00000000
--- a/upcaster/src/main/java/io/axoniq/dev/samples/api/PassengerSeatAdjustedEvent.java
+++ /dev/null
@@ -1,9 +0,0 @@
-package io.axoniq.dev.samples.api;
-
-public record PassengerSeatAdjustedEvent(
- String flightId,
- String passengerId,
- int seatNumber
-) {
-
-}
diff --git a/upcaster/src/main/java/io/axoniq/dev/samples/api/PassengerSeatsAdjustedEvent.java b/upcaster/src/main/java/io/axoniq/dev/samples/api/PassengerSeatsAdjustedEvent.java
deleted file mode 100644
index 2aa2154f..00000000
--- a/upcaster/src/main/java/io/axoniq/dev/samples/api/PassengerSeatsAdjustedEvent.java
+++ /dev/null
@@ -1,22 +0,0 @@
-package io.axoniq.dev.samples.api;
-
-import java.util.Map;
-
-/**
- * An old event of the flight domain that's now been deprecated.
- *
- * This event contained a complete list of all the passenger seat adjustments of a flight in one go. As time passed, the
- * application developers noticed they required an event for every separate change. Hence, they introduced that
- * {@link PassengerSeatAdjustedEvent} and added a one-to-many upcaster to make this adjustment.
- *
- * @author Steven van Beelen
- * @see io.axoniq.dev.samples.upcaster.json.PassengerSeatsToPassengerSeatAdjustedEventUpcaster
- * @deprecated in favor of singular {@link PassengerSeatAdjustedEvent}s
- */
-@Deprecated
-public record PassengerSeatsAdjustedEvent(
- String flightId,
- Map passengerSeats
-) {
-
-}
diff --git a/upcaster/src/main/java/io/axoniq/dev/samples/upcaster/json/EventUpcasterChainFactory.java b/upcaster/src/main/java/io/axoniq/dev/samples/upcaster/json/EventUpcasterChainFactory.java
deleted file mode 100644
index 8099101d..00000000
--- a/upcaster/src/main/java/io/axoniq/dev/samples/upcaster/json/EventUpcasterChainFactory.java
+++ /dev/null
@@ -1,53 +0,0 @@
-package io.axoniq.dev.samples.upcaster.json;
-
-import org.axonframework.config.Configurer;
-import org.axonframework.serialization.upcasting.event.EventUpcasterChain;
-
-import java.util.function.Function;
-
-/**
- * Utility class constructing the {@link EventUpcasterChain} to configure on the
- * {@link org.axonframework.eventsourcing.eventstore.EventStore}.
- *
- * In a Spring Boot environment exposing an {@code EventUpcasterChain} bean is sufficient for the framework to pick it
- * up correctly. To that end we can use the {@link #buildEventUpcasterChain()} method.
- *
- * When using Axon's {@link org.axonframework.config.Configurer} directly, you should configure all upcasters separately
- * by invoking the {@link org.axonframework.config.Configurer#registerEventUpcaster(Function)} method. The
- * {@link #configureUpcasters(Configurer)} shows how we should implement this.
- *
- * @author Yvonne Ceelie
- */
-@SuppressWarnings("unused")
-public abstract class EventUpcasterChainFactory {
-
- /**
- * Constructs an {@link EventUpcasterChain} combining all the upcasters of this application.
- *
- * Can be used to expose an {@code EventUpcasterChain} bean in a Spring environment.
- *
- * @return an {@link EventUpcasterChain}
- */
- public static EventUpcasterChain buildEventUpcasterChain() {
- return new EventUpcasterChain(
- new FlightDelayedEvent0_to_1Upcaster(),
- new PassengerSeatsToPassengerSeatAdjustedEventUpcaster()
- );
- }
-
- /**
- * Configures all the upcasters of this application with the given {@code configurer}.
- *
- * Can be utilized to configure upcasters if the application code uses Axon's {@link Configurer} directly.
- *
- * @param configurer the {@link Configurer} to register all upcasters with
- */
- public static void configureUpcasters(Configurer configurer) {
- configurer.registerEventUpcaster(config -> new FlightDelayedEvent0_to_1Upcaster())
- .registerEventUpcaster(config -> new PassengerSeatsToPassengerSeatAdjustedEventUpcaster());
- }
-
- private EventUpcasterChainFactory() {
- // Utility class
- }
-}
diff --git a/upcaster/src/main/java/io/axoniq/dev/samples/upcaster/json/FlightDelayedEvent0_to_1Upcaster.java b/upcaster/src/main/java/io/axoniq/dev/samples/upcaster/json/FlightDelayedEvent0_to_1Upcaster.java
deleted file mode 100644
index 6ba041e8..00000000
--- a/upcaster/src/main/java/io/axoniq/dev/samples/upcaster/json/FlightDelayedEvent0_to_1Upcaster.java
+++ /dev/null
@@ -1,68 +0,0 @@
-package io.axoniq.dev.samples.upcaster.json;
-
-import com.fasterxml.jackson.databind.JsonNode;
-import com.fasterxml.jackson.databind.node.JsonNodeFactory;
-import com.fasterxml.jackson.databind.node.ObjectNode;
-import io.axoniq.dev.samples.api.FlightDelayedEvent;
-import org.axonframework.serialization.SimpleSerializedType;
-import org.axonframework.serialization.upcasting.event.IntermediateEventRepresentation;
-import org.axonframework.serialization.upcasting.event.SingleEventUpcaster;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.lang.invoke.MethodHandles;
-
-/**
- * Upcaster upcasting the {@code FlightDelayedEvent} from revision {@code 0} to revision {@code 1}.
- *
- * This allows us to adjust the {@link FlightDelayedEvent} implementation to the new format. Part of the adjustment is
- * adding the {@link org.axonframework.serialization.Revision} annotation to the {@code FlightDelayedEvent}. The
- * annotation reflects the new version, defined as revision {@code 1}.
- *
- * The new format of the {@code FlightDelayedEvent} uses a {@link io.axoniq.dev.samples.api.Leg} object to contain the
- * {@code "origin"} and {@code "destination"} fields.
- *
- * @author Yvonne Ceelie
- */
-public class FlightDelayedEvent0_to_1Upcaster extends SingleEventUpcaster {
-
- private static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
-
- private static final String ORIGIN = "origin";
- private static final String DESTINATION = "destination";
- private static final String LEG = "leg";
-
- private final SimpleSerializedType sourceType =
- new SimpleSerializedType(FlightDelayedEvent.class.getTypeName(), null);
- private final SimpleSerializedType targetType =
- new SimpleSerializedType(FlightDelayedEvent.class.getTypeName(), "1.0");
-
- @Override
- protected boolean canUpcast(IntermediateEventRepresentation intermediateEventRepresentation) {
- return intermediateEventRepresentation.getType().equals(sourceType);
- }
-
- @Override
- protected IntermediateEventRepresentation doUpcast(
- IntermediateEventRepresentation intermediateEventRepresentation
- ) {
- logger.info("Upcast event: {}", intermediateEventRepresentation.getType());
- return intermediateEventRepresentation.upcastPayload(targetType, JsonNode.class, this::upcastEvent);
- }
-
- private JsonNode upcastEvent(JsonNode jsonNode) {
- if (!jsonNode.isObject()) {
- return jsonNode;
- }
- final ObjectNode root = (ObjectNode) jsonNode;
- // Create the leg node and set origin and destination in it
- ObjectNode leg = JsonNodeFactory.instance.objectNode();
- leg.set(ORIGIN, root.get(ORIGIN));
- leg.set(DESTINATION, root.get(DESTINATION));
- root.set(LEG, leg);
- // Remove origin and destination from the root
- root.remove(ORIGIN);
- root.remove(DESTINATION);
- return jsonNode;
- }
-}
diff --git a/upcaster/src/main/java/io/axoniq/dev/samples/upcaster/json/PassengerSeatsToPassengerSeatAdjustedEventUpcaster.java b/upcaster/src/main/java/io/axoniq/dev/samples/upcaster/json/PassengerSeatsToPassengerSeatAdjustedEventUpcaster.java
deleted file mode 100644
index ec9bf9cd..00000000
--- a/upcaster/src/main/java/io/axoniq/dev/samples/upcaster/json/PassengerSeatsToPassengerSeatAdjustedEventUpcaster.java
+++ /dev/null
@@ -1,81 +0,0 @@
-package io.axoniq.dev.samples.upcaster.json;
-
-import com.fasterxml.jackson.databind.JsonNode;
-import com.fasterxml.jackson.databind.node.ObjectNode;
-import com.fasterxml.jackson.databind.node.TextNode;
-import org.axonframework.serialization.SimpleSerializedType;
-import org.axonframework.serialization.upcasting.event.EventMultiUpcaster;
-import org.axonframework.serialization.upcasting.event.IntermediateEventRepresentation;
-
-import java.util.Map;
-import java.util.Spliterator;
-import java.util.Spliterators;
-import java.util.stream.Stream;
-import java.util.stream.StreamSupport;
-
-/**
- * An upcaster implementation that, instead of returning a single entry, returns a {@link Stream} of
- * {@link IntermediateEventRepresentation}s for the {@link io.axoniq.dev.samples.api.PassengerSeatsAdjustedEvent}.
- *
- * This {@link EventMultiUpcaster} implementation retrieves the {@code "passengerSeats"} collection from the deprecated
- * event. Doing os it is able to upcast the single event to the right amount of
- * {@link io.axoniq.dev.samples.api.PassengerSeatAdjustedEvent}s.
- *