This repository has been archived by the owner on Dec 3, 2021. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 5
/
HazelcastRateLimitRepository.java
102 lines (90 loc) · 4.05 KB
/
HazelcastRateLimitRepository.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
/**
* Copyright (C) 2015 The Gravitee team (http://gravitee.io)
*
* 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.gravitee.repository.hazelcast.ratelimit;
import com.hazelcast.core.HazelcastInstance;
import com.hazelcast.map.IMap;
import io.gravitee.repository.hazelcast.ratelimit.configuration.HazelcastRateLimitConfiguration;
import io.gravitee.repository.ratelimit.api.RateLimitRepository;
import io.gravitee.repository.ratelimit.model.RateLimit;
import io.reactivex.Completable;
import io.reactivex.Maybe;
import io.reactivex.Single;
import io.reactivex.SingleSource;
import io.reactivex.functions.Function;
import io.reactivex.schedulers.Schedulers;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import java.util.function.Supplier;
import javax.annotation.PostConstruct;
import org.springframework.beans.factory.annotation.Autowired;
/**
* @author David BRASSELY (david.brassely at graviteesource.com)
* @author GraviteeSource Team
*/
public class HazelcastRateLimitRepository implements RateLimitRepository<RateLimit> {
@Autowired
private HazelcastInstance hazelcastInstance;
@Autowired
private HazelcastRateLimitConfiguration configuration;
private IMap<String, RateLimit> counters;
@PostConstruct
public void afterPropertiesSet() {
counters = hazelcastInstance.getMap(configuration.getRateLimitMap());
}
@Override
public Single<RateLimit> incrementAndGet(String key, long weight, Supplier<RateLimit> supplier) {
final long now = System.currentTimeMillis();
Lock lock = hazelcastInstance.getCPSubsystem().getLock("lock-rl-" + key);
return Completable
.create(
emitter -> {
lock.lock();
emitter.onComplete();
}
)
.subscribeOn(Schedulers.computation())
.andThen(
Single.defer(
() ->
Maybe
.fromFuture(counters.getAsync(key).toCompletableFuture())
.switchIfEmpty((SingleSource<RateLimit>) observer -> observer.onSuccess(supplier.get()))
.flatMap(
(Function<RateLimit, SingleSource<RateLimit>>) rateLimit -> {
if (rateLimit.getResetTime() < now) {
rateLimit = supplier.get();
}
rateLimit.setCounter(rateLimit.getCounter() + weight);
final RateLimit finalRateLimit = rateLimit;
return Completable
.fromFuture(
counters
.setAsync(
rateLimit.getKey(),
rateLimit,
now - rateLimit.getResetTime(),
TimeUnit.MILLISECONDS
)
.toCompletableFuture()
)
.andThen(Single.defer(() -> Single.just(finalRateLimit)))
.doFinally(lock::unlock);
}
)
)
);
}
}