aboutsummaryrefslogtreecommitdiff
path: root/mavigator-uav/src/main/scala/mavigator/uav/Multiplexer.scala.disabled
diff options
context:
space:
mode:
Diffstat (limited to 'mavigator-uav/src/main/scala/mavigator/uav/Multiplexer.scala.disabled')
-rw-r--r--mavigator-uav/src/main/scala/mavigator/uav/Multiplexer.scala.disabled20
1 files changed, 20 insertions, 0 deletions
diff --git a/mavigator-uav/src/main/scala/mavigator/uav/Multiplexer.scala.disabled b/mavigator-uav/src/main/scala/mavigator/uav/Multiplexer.scala.disabled
new file mode 100644
index 0000000..5e48ea1
--- /dev/null
+++ b/mavigator-uav/src/main/scala/mavigator/uav/Multiplexer.scala.disabled
@@ -0,0 +1,20 @@
+package mavigator
+package uav
+
+import akka.stream.scaladsl._
+import akka.stream._
+import org.reactivestreams._
+
+private[uav] class Multiplexer[In, Out](service: Flow[In, Out, _])(implicit materializer: Materializer) {
+
+ private val endpoint: Flow[Out, In, (Publisher[Out], Subscriber[In])] = Flow.fromSinkAndSourceMat(
+ Sink.asPublisher[Out](fanout = true),
+ Source.asSubscriber[In])((pub, sub) => (pub, sub))
+
+ private lazy val (publisher, subscriber) = (service.joinMat(endpoint)(Keep.right)).run()
+
+ def connect(client: Flow[Out, In, _]) = {
+ Source.fromPublisher(publisher).via(client).to(Sink.ignore).run()
+ }
+
+}