Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions build.sbt
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import scalapb.compiler.Version.scalapbVersion

organization in ThisBuild := "beyondthelines"
version in ThisBuild := "0.0.7"
version in ThisBuild := "0.0.8"
bintrayOrganization in ThisBuild := Some(organization.value)
bintrayRepository in ThisBuild := "maven"
bintrayPackageLabels in ThisBuild := Seq("scala", "protobuf", "grpc", "monix")
Expand All @@ -15,7 +15,7 @@ lazy val runtime = (project in file("runtime"))
name := "GrpcMonixRuntime",
libraryDependencies ++= Seq(
"com.thesamet.scalapb" %% "scalapb-runtime-grpc" % scalapbVersion,
"io.monix" %% "monix" % "2.3.3"
"io.monix" %% "monix" % "3.1.0"
)
)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import com.google.protobuf.Descriptors._
import com.google.protobuf.ExtensionRegistry
import com.google.protobuf.compiler.PluginProtos.{CodeGeneratorRequest, CodeGeneratorResponse}
import scalapb.compiler.FunctionalPrinter.PrinterEndo
import scalapb.compiler.{DescriptorPimps, FunctionalPrinter, GeneratorParams, ProtobufGenerator, StreamType}
import scalapb.compiler.{DescriptorPimps, FunctionalPrinter, GeneratorParams, StreamType}
import scalapb.options.compiler.Scalapb

import scala.collection.JavaConverters._
Expand Down
33 changes: 26 additions & 7 deletions runtime/src/main/scala/grpcmonix/GrpcMonix.scala
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,11 @@ package grpcmonix

import com.google.common.util.concurrent.ListenableFuture
import io.grpc.stub.StreamObserver
import monix.eval.{Callback, Task}
import monix.eval.Task
import monix.execution.Ack.{Continue, Stop}
import monix.execution.{Ack, Scheduler}
import monix.execution.{Ack, Callback, Scheduler}
import monix.reactive.Observable
import monix.reactive.observables.ObservableLike.{Operator, Transformer}
import monix.reactive.Observable.Operator
import monix.reactive.observers.Subscriber
import monix.reactive.subjects.PublishSubject
import org.reactivestreams.{Subscriber => SubscriberR}
Expand All @@ -23,7 +23,7 @@ object GrpcMonix {
Grpc.guavaFuture2ScalaFuture(future)
}

def grpcOperatorToMonixOperator[I,O](grpcOperator: GrpcOperator[I,O]): Operator[I,O] = {
def grpcOperatorToMonixOperator[I, O](grpcOperator: GrpcOperator[I, O]): Operator[I, O] = {
outputSubsriber: Subscriber[O] =>
val outputObserver: StreamObserver[O] = monixSubscriberToGrpcObserver(outputSubsriber)
val inputObserver: StreamObserver[I] = grpcOperator(outputObserver)
Expand All @@ -33,22 +33,29 @@ object GrpcMonix {
def monixSubscriberToGrpcObserver[T](subscriber: Subscriber[T]): StreamObserver[T] =
new StreamObserver[T] {
override def onError(t: Throwable): Unit = subscriber.onError(t)

override def onCompleted(): Unit = subscriber.onComplete()

override def onNext(value: T): Unit = subscriber.onNext(value)
}

def reactiveSubscriberToGrpcObserver[T](subscriber: SubscriberR[_ >: T]): StreamObserver[T] =
new StreamObserver[T] {
override def onError(t: Throwable): Unit = subscriber.onError(t)

override def onCompleted(): Unit = subscriber.onComplete()

override def onNext(value: T): Unit = subscriber.onNext(value)
}

def grpcObserverToMonixSubscriber[T](observer: StreamObserver[T], s: Scheduler): Subscriber[T] =
new Subscriber[T] {
override implicit def scheduler: Scheduler = s

override def onError(t: Throwable): Unit = observer.onError(t)

override def onComplete(): Unit = observer.onCompleted()

override def onNext(value: T): Future[Ack] =
try {
observer.onNext(value)
Expand All @@ -60,10 +67,11 @@ object GrpcMonix {
}
}

def grpcObserverToMonixCallback[T](observer: StreamObserver[T]): Callback[T] =
new Callback[T] {
def grpcObserverToMonixCallback[A](observer: StreamObserver[A]): Callback[Throwable, A] =
new Callback[Throwable, A] {
override def onError(t: Throwable): Unit = observer.onError(t)
override def onSuccess(value: T): Unit = {

override def onSuccess(value: A): Unit = {
observer.onNext(value)
observer.onCompleted()
}
Expand All @@ -80,9 +88,20 @@ object GrpcMonix {
subject.transform(transformer).subscribe(subscriber)

override implicit def scheduler: Scheduler = subscriber.scheduler

override def onError(t: Throwable): Unit = subject.onError(t)

override def onComplete(): Unit = subject.onComplete()

override def onNext(value: I): Future[Ack] = subject.onNext(value)
}

type Transformer[A, B] = Observable[A] => Observable[B]

implicit class ObservableExtension[A](ob: Observable[A]) {
def transform[B](transformer: Transformer[A, B]): Observable[B] = {
transformer(ob)
}
}

}