-
Notifications
You must be signed in to change notification settings - Fork 215
/
MongoThingsSearchUpdaterPersistence.java
131 lines (115 loc) · 6.04 KB
/
MongoThingsSearchUpdaterPersistence.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
/*
* Copyright (c) 2017 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.thingsearch.service.persistence.write.impl;
import static com.mongodb.client.model.Filters.elemMatch;
import static com.mongodb.client.model.Filters.in;
import static com.mongodb.client.model.Filters.or;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
import org.bson.BsonDateTime;
import org.bson.BsonDocument;
import org.bson.BsonInt32;
import org.bson.BsonString;
import org.bson.Document;
import org.bson.conversions.Bson;
import org.eclipse.ditto.internal.models.streaming.EntityIdWithRevision;
import org.eclipse.ditto.policies.api.PolicyTag;
import org.eclipse.ditto.policies.model.PolicyId;
import org.eclipse.ditto.things.model.ThingId;
import org.eclipse.ditto.thingsearch.api.PolicyReferenceTag;
import org.eclipse.ditto.thingsearch.service.common.config.SearchPersistenceConfig;
import org.eclipse.ditto.thingsearch.service.persistence.PersistenceConstants;
import org.eclipse.ditto.thingsearch.service.persistence.write.ThingsSearchUpdaterPersistence;
import org.eclipse.ditto.thingsearch.service.persistence.write.model.AbstractWriteModel;
import org.reactivestreams.Publisher;
import com.mongodb.client.model.UpdateManyModel;
import com.mongodb.client.model.UpdateOptions;
import com.mongodb.client.model.WriteModel;
import com.mongodb.reactivestreams.client.MongoCollection;
import com.mongodb.reactivestreams.client.MongoDatabase;
import akka.NotUsed;
import akka.japi.pf.PFBuilder;
import akka.stream.javadsl.Source;
/**
* MongoDB specific implementation of the {@link org.eclipse.ditto.thingsearch.service.persistence.write.ThingsSearchUpdaterPersistence}.
*/
public final class MongoThingsSearchUpdaterPersistence implements ThingsSearchUpdaterPersistence {
private final MongoCollection<Document> collection;
private MongoThingsSearchUpdaterPersistence(final MongoDatabase database,
final SearchPersistenceConfig updaterPersistenceConfig) {
collection = database.getCollection(PersistenceConstants.THINGS_COLLECTION_NAME)
.withReadConcern(updaterPersistenceConfig.readConcern().getMongoReadConcern())
.withReadPreference(updaterPersistenceConfig.readPreference().getMongoReadPreference());
}
/**
* Constructor.
*
* @param database the database.
* @param updaterPersistenceConfig the updater persistence config to use.
*/
public static ThingsSearchUpdaterPersistence of(final MongoDatabase database,
final SearchPersistenceConfig updaterPersistenceConfig) {
return new MongoThingsSearchUpdaterPersistence(database, updaterPersistenceConfig);
}
@Override
public Source<PolicyReferenceTag, NotUsed> getPolicyReferenceTags(final Map<PolicyId, Long> policyRevisions) {
final Bson usedAsThingPolicy = in(PersistenceConstants.FIELD_POLICY_ID, policyRevisions.keySet()
.stream()
.map(String::valueOf)
.collect(Collectors.toSet()));
final Bson isReferencedPolicy = elemMatch(PersistenceConstants.FIELD_REFERENCED_POLICIES,
in(EntityIdWithRevision.JsonFields.ENTITY_ID.getPointer().toString(), policyRevisions.keySet()
.stream()
.map(String::valueOf)
.collect(Collectors.toSet())));
final Bson filter = or(
usedAsThingPolicy, // This is only required for backwards compatibility.
isReferencedPolicy
);
final Publisher<Document> publisher =
collection.find(filter).projection(new Document()
.append(PersistenceConstants.FIELD_ID, new BsonInt32(1))
.append(PersistenceConstants.FIELD_POLICY_ID, new BsonInt32(1)));
return Source.fromPublisher(publisher)
.mapConcat(doc -> {
final ThingId thingId = ThingId.of(doc.getString(PersistenceConstants.FIELD_ID));
final String policyIdString = doc.getString(PersistenceConstants.FIELD_POLICY_ID);
final PolicyId policyId = PolicyId.of(policyIdString);
final Long revision = policyRevisions.get(policyId);
if (revision == null) {
return Collections.emptyList();
} else {
final PolicyTag policyTag = PolicyTag.of(policyId, revision);
return Collections.singletonList(PolicyReferenceTag.of(thingId, policyTag));
}
});
}
@Override
public Source<List<Throwable>, NotUsed> purge(final CharSequence namespace) {
final Bson filter = thingNamespaceFilter(namespace);
final Bson update = new BsonDocument().append(AbstractWriteModel.SET,
new BsonDocument().append(PersistenceConstants.FIELD_DELETE_AT, new BsonDateTime(0L)));
final UpdateOptions updateOptions = new UpdateOptions().bypassDocumentValidation(true);
final WriteModel<Document> writeModel = new UpdateManyModel<>(filter, update, updateOptions);
return Source.fromPublisher(collection.bulkWrite(Collections.singletonList(writeModel)))
.map(bulkWriteResult -> Collections.<Throwable>emptyList())
.recoverWithRetries(1, new PFBuilder<Throwable, Source<List<Throwable>, NotUsed>>()
.matchAny(throwable -> Source.single(Collections.singletonList(throwable)))
.build());
}
private Document thingNamespaceFilter(final CharSequence namespace) {
return new Document().append(PersistenceConstants.FIELD_NAMESPACE, new BsonString(namespace.toString()));
}
}