-
Notifications
You must be signed in to change notification settings - Fork 214
/
LiveSignalPubSubFactory.java
89 lines (76 loc) · 3.66 KB
/
LiveSignalPubSubFactory.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
/*
* Copyright (c) 2022 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.internal.utils.pubsubthings;
import java.util.Collection;
import java.util.Collections;
import org.eclipse.ditto.base.model.signals.Signal;
import org.eclipse.ditto.base.model.signals.SignalWithEntityId;
import org.eclipse.ditto.internal.utils.pubsub.AbstractPubSubFactory;
import org.eclipse.ditto.internal.utils.pubsub.DistributedAcks;
import org.eclipse.ditto.internal.utils.pubsub.StreamingType;
import org.eclipse.ditto.internal.utils.pubsub.extractors.AckExtractor;
import org.eclipse.ditto.internal.utils.pubsub.extractors.PubSubTopicExtractor;
import org.eclipse.ditto.internal.utils.pubsub.extractors.ReadSubjectExtractor;
import org.eclipse.ditto.things.model.ThingId;
import org.apache.pekko.actor.ActorContext;
import org.apache.pekko.actor.ActorRefFactory;
import org.apache.pekko.actor.ActorSystem;
/**
* Pub-sub factory for live signals.
*/
final class LiveSignalPubSubFactory extends AbstractPubSubFactory<SignalWithEntityId<?>> {
private static final AckExtractor<SignalWithEntityId<?>> ACK_EXTRACTOR =
AckExtractor.of(LiveSignalPubSubFactory::getThingId, Signal::getDittoHeaders);
private static final DDataProvider PROVIDER = DDataProvider.of("live-signal-aware");
@SuppressWarnings("unchecked")
private LiveSignalPubSubFactory(final ActorRefFactory actorRefFactory,
final ActorSystem actorSystem,
final PubSubTopicExtractor<SignalWithEntityId<?>> topicExtractor,
final DistributedAcks distributedAcks) {
super(actorRefFactory, actorSystem, (Class<SignalWithEntityId<?>>) (Object) Signal.class, topicExtractor, PROVIDER,
ACK_EXTRACTOR, distributedAcks);
}
/**
* Create a pubsub factory for live signals from an actor context.
*
* @param context context of the actor under which the publisher and subscriber actors are started.
* @param distributedAcks the distributed acks interface.
* @return the thing
*/
public static LiveSignalPubSubFactory of(final ActorContext context, final DistributedAcks distributedAcks) {
return new LiveSignalPubSubFactory(context, context.system(), topicExtractor(), distributedAcks);
}
/**
* Create a pubsub factory for live signals from an actor system.
*
* @param system the actor system.
* @param distributedAcks the distributed acks interface.
* @return the thing
*/
public static LiveSignalPubSubFactory of(final ActorSystem system, final DistributedAcks distributedAcks) {
return new LiveSignalPubSubFactory(system, system, topicExtractor(), distributedAcks);
}
private static Collection<String> getStreamingTypeTopic(final Signal<?> signal) {
return StreamingType.fromSignal(signal)
.map(StreamingType::getDistributedPubSubTopic)
.map(Collections::singleton)
.orElse(Collections.emptySet());
}
private static PubSubTopicExtractor<SignalWithEntityId<?>> topicExtractor() {
return ReadSubjectExtractor.<SignalWithEntityId<?>>of().with(LiveSignalPubSubFactory::getStreamingTypeTopic);
}
// precondition: all live signals are thing signals.
private static ThingId getThingId(final SignalWithEntityId<?> signal) {
return ThingId.of(signal.getEntityId());
}
}