-
Notifications
You must be signed in to change notification settings - Fork 47
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Added ActorConfig based support to specify required message metadata …
…generators
- Loading branch information
1 parent
9f18ad5
commit ebf0c38
Showing
21 changed files
with
750 additions
and
67 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
26 changes: 23 additions & 3 deletions
26
src/main/java/io/appform/dropwizard/actors/actor/MessageMetadata.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,15 +1,35 @@ | ||
package io.appform.dropwizard.actors.actor; | ||
|
||
import com.rabbitmq.client.AMQP; | ||
import java.util.Date; | ||
import java.util.Map; | ||
import lombok.AllArgsConstructor; | ||
import lombok.Builder; | ||
import lombok.Data; | ||
import lombok.NoArgsConstructor; | ||
|
||
@Data | ||
@Builder | ||
@NoArgsConstructor | ||
@AllArgsConstructor | ||
public final class MessageMetadata { | ||
|
||
private String contentType; | ||
private String contentEncoding; | ||
private Map<String,Object> headers; | ||
private Integer deliveryMode; | ||
private Integer priority; | ||
private String correlationId; | ||
private String replyTo; | ||
private String expiration; | ||
private String messageId; | ||
private Date msgSentTimestamp; | ||
private String type; | ||
private String userId; | ||
private String appId; | ||
private String clusterId; | ||
|
||
// custom fields | ||
private boolean redelivered; | ||
private long delayInMs; | ||
private AMQP.BasicProperties properties; | ||
|
||
private boolean expired; | ||
} |
30 changes: 30 additions & 0 deletions
30
src/main/java/io/appform/dropwizard/actors/actor/metadata/MessageMetaContext.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,30 @@ | ||
package io.appform.dropwizard.actors.actor.metadata; | ||
|
||
import java.util.Date; | ||
import java.util.Map; | ||
import lombok.Builder; | ||
import lombok.Data; | ||
|
||
/** | ||
* Ref : <a href="https://www.rabbitmq.com/amqp-0-9-1-reference.html#class.basic">AMQP properties</a> | ||
*/ | ||
@Data | ||
@Builder | ||
public class MessageMetaContext { | ||
|
||
private boolean redelivered; | ||
private String contentType; | ||
private String contentEncoding; | ||
private Map<String,Object> headers; | ||
private Integer deliveryMode; | ||
private Integer priority; | ||
private String correlationId; | ||
private String replyTo; | ||
private String expiration; | ||
private String messageId; | ||
private Date timestamp; | ||
private String type; | ||
private String userId; | ||
private String appId; | ||
private String clusterId; | ||
} |
67 changes: 67 additions & 0 deletions
67
src/main/java/io/appform/dropwizard/actors/actor/metadata/MessageMetadataProvider.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,67 @@ | ||
package io.appform.dropwizard.actors.actor.metadata; | ||
|
||
import com.rabbitmq.client.AMQP; | ||
import com.rabbitmq.client.Envelope; | ||
import io.appform.dropwizard.actors.actor.MessageMetadata; | ||
import io.appform.dropwizard.actors.actor.metadata.generators.MessageMetadataGenerator; | ||
import io.appform.dropwizard.actors.common.RabbitmqActorException; | ||
import java.lang.reflect.InvocationTargetException; | ||
import java.util.List; | ||
import java.util.Objects; | ||
import lombok.val; | ||
|
||
public class MessageMetadataProvider { | ||
|
||
private final List<MessageMetadataGenerator> messageMetadataGenerators; | ||
|
||
public MessageMetadataProvider(List<String> messageMetaGeneratorClasses) { | ||
messageMetadataGenerators = configureMessageMetadataGenerators(messageMetaGeneratorClasses); | ||
} | ||
|
||
public MessageMetadata createMetadata(final Envelope envelope, final AMQP.BasicProperties properties) { | ||
Objects.requireNonNull(envelope, "Envelope should not be null"); | ||
Objects.requireNonNull(properties, "AMQP.BasicProperties should not be null"); | ||
val messageMetaContext = MessageMetaContext.builder() | ||
.redelivered(envelope.isRedeliver()) | ||
.contentType(properties.getContentType()) | ||
.contentEncoding(properties.getContentEncoding()) | ||
.headers(properties.getHeaders()) | ||
.deliveryMode(properties.getDeliveryMode()) | ||
.priority(properties.getPriority()) | ||
.correlationId(properties.getCorrelationId()) | ||
.replyTo(properties.getReplyTo()) | ||
.expiration(properties.getExpiration()) | ||
.messageId(properties.getMessageId()) | ||
.timestamp(properties.getTimestamp()) | ||
.type(properties.getType()) | ||
.userId(properties.getUserId()) | ||
.appId(properties.getAppId()) | ||
.clusterId(properties.getClusterId()) | ||
.build(); | ||
|
||
val messageMetadata = MessageMetadata.builder().build(); | ||
messageMetadataGenerators.forEach(generator -> generator.generate(messageMetaContext, messageMetadata)); | ||
return messageMetadata; | ||
} | ||
|
||
private List<MessageMetadataGenerator> configureMessageMetadataGenerators(List<String> messageMetaGeneratorClasses) { | ||
if (null == messageMetaGeneratorClasses) { | ||
return List.of(); | ||
} | ||
|
||
return messageMetaGeneratorClasses.stream() | ||
.distinct() | ||
.map(className -> { | ||
try { | ||
return (MessageMetadataGenerator) Class.forName(className) | ||
.getDeclaredConstructor() | ||
.newInstance(); | ||
} catch (ClassNotFoundException | NoSuchMethodException | InstantiationException | | ||
IllegalAccessException | InvocationTargetException e) { | ||
throw RabbitmqActorException.propagate("Failed to initialize messageMetadataGenerator " | ||
+ className, e); | ||
} | ||
}) | ||
.toList(); | ||
} | ||
} |
23 changes: 23 additions & 0 deletions
23
...ava/io/appform/dropwizard/actors/actor/metadata/generators/MessageDelayMetaGenerator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,23 @@ | ||
package io.appform.dropwizard.actors.actor.metadata.generators; | ||
|
||
import static io.appform.dropwizard.actors.common.Constants.MESSAGE_PUBLISHED_TEXT; | ||
|
||
import io.appform.dropwizard.actors.actor.MessageMetadata; | ||
import io.appform.dropwizard.actors.actor.metadata.MessageMetaContext; | ||
import io.appform.dropwizard.actors.utils.CommonUtils; | ||
import java.time.Instant; | ||
|
||
public class MessageDelayMetaGenerator implements MessageMetadataGenerator { | ||
|
||
@Override | ||
public void generate(final MessageMetaContext messageMetaContext, MessageMetadata messageMetadata) { | ||
messageMetadata.setDelayInMs(calculateDelayInMs(messageMetaContext)); | ||
} | ||
|
||
private long calculateDelayInMs(MessageMetaContext messageMetaContext) { | ||
return CommonUtils.extractMessagePropertiesHeader(messageMetaContext.getHeaders(), | ||
MESSAGE_PUBLISHED_TEXT, Long.class) | ||
.map(publishedAt -> Math.max(Instant.now().toEpochMilli() - publishedAt, 0)) | ||
.orElse(0L); | ||
} | ||
} |
24 changes: 24 additions & 0 deletions
24
...a/io/appform/dropwizard/actors/actor/metadata/generators/MessageExpiredMetaGenerator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,24 @@ | ||
package io.appform.dropwizard.actors.actor.metadata.generators; | ||
|
||
import static io.appform.dropwizard.actors.common.Constants.MESSAGE_EXPIRY_TEXT; | ||
|
||
import io.appform.dropwizard.actors.actor.MessageMetadata; | ||
import io.appform.dropwizard.actors.actor.metadata.MessageMetaContext; | ||
import io.appform.dropwizard.actors.utils.CommonUtils; | ||
import java.time.Instant; | ||
|
||
public class MessageExpiredMetaGenerator implements MessageMetadataGenerator { | ||
|
||
@Override | ||
public void generate(MessageMetaContext messageMetaContext, MessageMetadata messageMetadata) { | ||
messageMetadata.setExpired(isExpired(messageMetaContext)); | ||
} | ||
|
||
private boolean isExpired(MessageMetaContext messageMetaContext) { | ||
return CommonUtils.extractMessagePropertiesHeader(messageMetaContext.getHeaders(), | ||
MESSAGE_EXPIRY_TEXT, Long.class) | ||
.filter(expiry -> Instant.now().toEpochMilli() >= expiry) | ||
.isPresent(); | ||
} | ||
|
||
} |
13 changes: 13 additions & 0 deletions
13
...a/io/appform/dropwizard/actors/actor/metadata/generators/MessageHeadersMetaGenerator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,13 @@ | ||
package io.appform.dropwizard.actors.actor.metadata.generators; | ||
|
||
import io.appform.dropwizard.actors.actor.MessageMetadata; | ||
import io.appform.dropwizard.actors.actor.metadata.MessageMetaContext; | ||
|
||
public class MessageHeadersMetaGenerator implements MessageMetadataGenerator { | ||
|
||
@Override | ||
public void generate(final MessageMetaContext messageMetaContext, MessageMetadata messageMetadata) { | ||
messageMetadata.setHeaders(messageMetaContext.getHeaders()); | ||
} | ||
|
||
} |
25 changes: 25 additions & 0 deletions
25
...java/io/appform/dropwizard/actors/actor/metadata/generators/MessageMetaInfoGenerator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,25 @@ | ||
package io.appform.dropwizard.actors.actor.metadata.generators; | ||
|
||
import io.appform.dropwizard.actors.actor.MessageMetadata; | ||
import io.appform.dropwizard.actors.actor.metadata.MessageMetaContext; | ||
|
||
public class MessageMetaInfoGenerator implements MessageMetadataGenerator { | ||
|
||
@Override | ||
public void generate(final MessageMetaContext messageMetaContext, MessageMetadata messageMetadata) { | ||
messageMetadata.setContentType(messageMetaContext.getContentType()); | ||
messageMetadata.setContentEncoding(messageMetaContext.getContentEncoding()); | ||
messageMetadata.setDeliveryMode(messageMetaContext.getDeliveryMode()); | ||
messageMetadata.setPriority(messageMetaContext.getPriority()); | ||
messageMetadata.setCorrelationId(messageMetaContext.getCorrelationId()); | ||
messageMetadata.setReplyTo(messageMetaContext.getReplyTo()); | ||
messageMetadata.setExpiration(messageMetaContext.getExpiration()); | ||
messageMetadata.setMessageId(messageMetaContext.getMessageId()); | ||
messageMetadata.setMsgSentTimestamp(messageMetaContext.getTimestamp()); | ||
messageMetadata.setType(messageMetaContext.getType()); | ||
messageMetadata.setUserId(messageMetaContext.getUserId()); | ||
messageMetadata.setAppId(messageMetaContext.getAppId()); | ||
messageMetadata.setClusterId(messageMetaContext.getClusterId()); | ||
} | ||
|
||
} |
11 changes: 11 additions & 0 deletions
11
...java/io/appform/dropwizard/actors/actor/metadata/generators/MessageMetadataGenerator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,11 @@ | ||
package io.appform.dropwizard.actors.actor.metadata.generators; | ||
|
||
import io.appform.dropwizard.actors.actor.MessageMetadata; | ||
import io.appform.dropwizard.actors.actor.metadata.MessageMetaContext; | ||
|
||
@FunctionalInterface | ||
public interface MessageMetadataGenerator { | ||
|
||
void generate(final MessageMetaContext messageMetaContext, final MessageMetadata messageMetadata); | ||
|
||
} |
13 changes: 13 additions & 0 deletions
13
.../appform/dropwizard/actors/actor/metadata/generators/MessageReDeliveredMetaGenerator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,13 @@ | ||
package io.appform.dropwizard.actors.actor.metadata.generators; | ||
|
||
import io.appform.dropwizard.actors.actor.MessageMetadata; | ||
import io.appform.dropwizard.actors.actor.metadata.MessageMetaContext; | ||
|
||
public class MessageReDeliveredMetaGenerator implements MessageMetadataGenerator { | ||
|
||
@Override | ||
public void generate(final MessageMetaContext messageMetaContext, MessageMetadata messageMetadata) { | ||
messageMetadata.setRedelivered(messageMetaContext.isRedelivered()); | ||
} | ||
|
||
} |
Oops, something went wrong.