/
AbstractEnforcerActor.java
156 lines (134 loc) · 6.47 KB
/
AbstractEnforcerActor.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
/*
* Copyright (c) 2019 Contributors to the Eclipse Foundation
*
* See the NOTICE file(s) distributed with this work for additional
* information regarding copyright ownership.
*
* This program and the accompanying materials are made available under the
* terms of the Eclipse Public License 2.0 which is available at
* http://www.eclipse.org/legal/epl-2.0
*
* SPDX-License-Identifier: EPL-2.0
*/
package org.eclipse.ditto.services.concierge.enforcement;
import java.time.Duration;
import java.util.concurrent.Executor;
import java.util.concurrent.TimeUnit;
import javax.annotation.Nullable;
import org.eclipse.ditto.model.base.headers.WithDittoHeaders;
import org.eclipse.ditto.model.enforcers.Enforcer;
import org.eclipse.ditto.services.models.concierge.EntityId;
import org.eclipse.ditto.services.models.concierge.InvalidateCacheEntry;
import org.eclipse.ditto.services.utils.akka.controlflow.AbstractGraphActor;
import org.eclipse.ditto.services.utils.cache.Cache;
import org.eclipse.ditto.services.utils.cache.CaffeineCache;
import org.eclipse.ditto.services.utils.cache.entry.Entry;
import org.eclipse.ditto.services.utils.metrics.DittoMetrics;
import org.eclipse.ditto.services.utils.metrics.instruments.timer.ExpiringTimerBuilder;
import org.eclipse.ditto.services.utils.metrics.instruments.timer.StartedTimer;
import org.eclipse.ditto.signals.base.Signal;
import org.eclipse.ditto.signals.commands.base.Command;
import com.github.benmanes.caffeine.cache.Caffeine;
import akka.NotUsed;
import akka.actor.ActorRef;
import akka.cluster.pubsub.DistributedPubSubMediator;
import akka.japi.pf.ReceiveBuilder;
import akka.stream.javadsl.Flow;
/**
* Extensible actor to execute enforcement behavior.
*/
public abstract class AbstractEnforcerActor extends AbstractGraphActor<Contextual<WithDittoHeaders>> {
private static final String TIMER_NAME = "concierge_enforcements";
/**
* Contextual information about this actor.
*/
protected final Contextual<WithDittoHeaders> contextual;
@Nullable
private final Cache<EntityId, Entry<EntityId>> thingIdCache;
@Nullable
private final Cache<EntityId, Entry<Enforcer>> aclEnforcerCache;
@Nullable
private final Cache<EntityId, Entry<Enforcer>> policyEnforcerCache;
/**
* Create an instance of this actor.
*
* @param pubSubMediator Akka pub-sub-mediator.
* @param conciergeForwarder the concierge forwarder.
* @param enforcerExecutor executor for enforcement steps.
* @param askTimeout how long to wait for entity actors.
* @param bufferSize the buffer size used for the Source queue.
* @param parallelism parallelism to use for processing messages in parallel.
* @param thingIdCache the cache for Thing IDs to either ACL or Policy ID.
* @param aclEnforcerCache the ACL cache.
* @param policyEnforcerCache the Policy cache.
*/
protected AbstractEnforcerActor(final ActorRef pubSubMediator,
final ActorRef conciergeForwarder,
final Executor enforcerExecutor,
final Duration askTimeout,
final int bufferSize,
final int parallelism,
@Nullable final Cache<EntityId, Entry<EntityId>> thingIdCache,
@Nullable final Cache<EntityId, Entry<Enforcer>> aclEnforcerCache,
@Nullable final Cache<EntityId, Entry<Enforcer>> policyEnforcerCache) {
super(bufferSize, parallelism);
this.thingIdCache = thingIdCache;
this.aclEnforcerCache = aclEnforcerCache;
this.policyEnforcerCache = policyEnforcerCache;
contextual = new Contextual<>(null, getSelf(), getContext().getSystem().deadLetters(),
pubSubMediator, conciergeForwarder, enforcerExecutor, askTimeout, log, null, null,
createResponseReceiversCache());
// register for sending messages via pub/sub to this enforcer
// used for receiving cache invalidations from brother concierge nodes
pubSubMediator.tell(new DistributedPubSubMediator.Put(getSelf()), getSelf());
}
@Override
protected void preEnhancement(final ReceiveBuilder receiveBuilder) {
receiveBuilder.match(InvalidateCacheEntry.class, invalidateCacheEntry -> {
log.debug("received <{}>", invalidateCacheEntry);
final EntityId entityId = invalidateCacheEntry.getEntityId();
invalidateCaches(entityId);
});
}
private void invalidateCaches(final EntityId entityId) {
if (thingIdCache != null) {
final boolean invalidated = thingIdCache.invalidate(entityId);
log.debug("thingId cache for entity id <{}> was invalidated: {}", entityId, invalidated);
}
if (aclEnforcerCache != null) {
final boolean invalidated = aclEnforcerCache.invalidate(entityId);
log.debug("acl enforcer cache for entity id <{}> was invalidated: {}", entityId, invalidated);
}
if (policyEnforcerCache != null) {
final boolean invalidated = policyEnforcerCache.invalidate(entityId);
log.debug("policy enforcer cache for entity id <{}> was invalidated: {}", entityId, invalidated);
}
}
@Override
protected Contextual<WithDittoHeaders> beforeHandleMessage(final Contextual<WithDittoHeaders> contextual) {
final StartedTimer timer = createTimer(contextual.getMessage());
return contextual.withTimer(timer);
}
private StartedTimer createTimer(final WithDittoHeaders withDittoHeaders) {
final ExpiringTimerBuilder timerBuilder = DittoMetrics.expiringTimer(TIMER_NAME);
withDittoHeaders.getDittoHeaders().getChannel().ifPresent(channel ->
timerBuilder.tag("channel", channel)
);
if (withDittoHeaders instanceof Signal) {
timerBuilder.tag("resource", ((Signal) withDittoHeaders).getResourceType());
}
if (withDittoHeaders instanceof Command) {
timerBuilder.tag("category", ((Command) withDittoHeaders).getCategory().name().toLowerCase());
}
return timerBuilder.build();
}
@Override
protected abstract Flow<Contextual<WithDittoHeaders>, Contextual<WithDittoHeaders>, NotUsed> getHandler();
@Override
protected Contextual<WithDittoHeaders> mapMessage(final WithDittoHeaders message) {
return contextual.withReceivedMessage(message, getSender());
}
private static Cache<String, ActorRef> createResponseReceiversCache() {
return CaffeineCache.of(Caffeine.newBuilder().expireAfterWrite(120, TimeUnit.SECONDS));
}
}