Last active
July 28, 2022 13:38
-
-
Save irajhedayati/b923179ba7ac4abd27daa480882e1584 to your computer and use it in GitHub Desktop.
A Pub/Sub subscriber wrapper that supports local emulator
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
class SubscriberWrapper(subscriber: Subscriber) extends AutoCloseable { | |
override def close(): Unit = subscriber.stopAsync() | |
def startAsync: ApiService = subscriber.startAsync() | |
def awaitTerminated(): Unit = subscriber.awaitTerminated() | |
} | |
object SubscriberWrapper { | |
def apply(projectId: String, subscription: String, localPubSubEndpoint: Option[String], receiver: MessageReceiver): SubscriberWrapper = { | |
val baseBuilder = Subscriber.newBuilder(ProjectSubscriptionName.format(projectId, subscription), receiver) | |
val subscriber = | |
localPubSubEndpoint | |
.map(h => handleLocalExecution(baseBuilder, h)) | |
.map(_.build()) | |
.getOrElse(baseBuilder.build()) | |
new SubscriberWrapper(subscriber) | |
} | |
private def handleLocalExecution(baseBuilder: Subscriber.Builder, localPubSubEndpoint: String): Subscriber.Builder = { | |
val channel = ManagedChannelBuilder | |
.forTarget(localPubSubEndpoint) | |
.usePlaintext() | |
.asInstanceOf[ManagedChannelBuilder[NettyChannelBuilder]] | |
.build() | |
baseBuilder | |
.setChannelProvider(FixedTransportChannelProvider.create(GrpcTransportChannel.create(channel))) | |
.setCredentialsProvider(NoCredentialsProvider.create) | |
} | |
} |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment