Metric: Processing time
[dcaegen2/collectors/hv-ves.git] / sources / hv-collector-ct / src / test / kotlin / org / onap / dcae / collectors / veshv / tests / fakes / sink.kt
1 /*
2  * ============LICENSE_START=======================================================
3  * dcaegen2-collectors-veshv
4  * ================================================================================
5  * Copyright (C) 2018 NOKIA
6  * ================================================================================
7  * Licensed under the Apache License, Version 2.0 (the "License");
8  * you may not use this file except in compliance with the License.
9  * You may obtain a copy of the License at
10  *
11  *      http://www.apache.org/licenses/LICENSE-2.0
12  *
13  * Unless required by applicable law or agreed to in writing, software
14  * distributed under the License is distributed on an "AS IS" BASIS,
15  * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
16  * See the License for the specific language governing permissions and
17  * limitations under the License.
18  * ============LICENSE_END=========================================================
19  */
20 package org.onap.dcae.collectors.veshv.tests.fakes
21
22 import arrow.core.identity
23 import org.onap.dcae.collectors.veshv.boundary.Sink
24 import org.onap.dcae.collectors.veshv.model.RoutedMessage
25 import org.reactivestreams.Publisher
26 import reactor.core.publisher.Flux
27 import java.util.*
28 import java.util.concurrent.ConcurrentLinkedDeque
29 import java.util.concurrent.atomic.AtomicLong
30 import java.util.function.Function
31
32 /**
33  * @author Piotr Jaszczyk <piotr.jaszczyk@nokia.com>
34  * @since May 2018
35  */
36 class StoringSink : Sink {
37     private val sent: Deque<RoutedMessage> = ConcurrentLinkedDeque()
38
39     val sentMessages: List<RoutedMessage>
40         get() = sent.toList()
41
42     override fun send(messages: Flux<RoutedMessage>): Flux<RoutedMessage> {
43         return messages.doOnNext(sent::addLast)
44     }
45 }
46
47 /**
48  * @author Piotr Jaszczyk <piotr.jaszczyk@nokia.com>
49  * @since May 2018
50  */
51 class CountingSink : Sink {
52     private val atomicCount = AtomicLong(0)
53
54     val count: Long
55         get() = atomicCount.get()
56
57     override fun send(messages: Flux<RoutedMessage>): Flux<RoutedMessage> {
58         return messages.doOnNext {
59             atomicCount.incrementAndGet()
60         }
61     }
62 }
63
64
65 open class ProcessingSink(val transformer: (Flux<RoutedMessage>) -> Publisher<RoutedMessage>) : Sink {
66     override fun send(messages: Flux<RoutedMessage>): Flux<RoutedMessage> = messages.transform(transformer)
67 }
68
69 class NoOpSink : ProcessingSink(::identity)