Reconnect text client streams (#8429)

This commit is contained in:
2026-07-10 22:16:18 -07:00
committed by GitHub
parent 6a0162109b
commit 90d1cf702f
@@ -119,25 +119,26 @@ final class EagleTextClient(config: Config) {
baseStub.withInterceptors(MetadataUtils.newAttachHeadersInterceptor(metadata))
}
private val state = new ClientState
private val running = new AtomicBoolean(true)
private val responseObserver = new StreamObserver[UpdateStreamResponse] {
private val state = new ClientState
private val running = new AtomicBoolean(true)
private val reconnecting = new AtomicBoolean(false)
private var requestObserver: Option[StreamObserver[UpdateStreamRequest]] = None
private val responseObserver: StreamObserver[UpdateStreamResponse] = new StreamObserver[UpdateStreamResponse] {
override def onNext(response: UpdateStreamResponse): Unit =
state.handle(response)
override def onError(t: Throwable): Unit = {
println(s"stream error: ${t.getMessage}")
running.set(false)
reconnect()
}
override def onCompleted(): Unit = {
println("stream completed")
running.set(false)
reconnect()
}
}
private val requestObserver: StreamObserver[UpdateStreamRequest] =
stub.streamUpdates(responseObserver)
openRequestStream()
def run(): Unit =
try {
@@ -147,7 +148,7 @@ final class EagleTextClient(config: Config) {
println("Type 'help' for commands.")
repl()
} finally {
val _ = Try(requestObserver.onCompleted())
requestObserver.foreach(observer => Try(observer.onCompleted()))
channel.shutdown()
val _ = channel.awaitTermination(2, TimeUnit.SECONDS)
}
@@ -405,9 +406,41 @@ final class EagleTextClient(config: Config) {
}
private def send(request: UpdateStreamRequest): Unit =
Try(requestObserver.onNext(request)) match {
case Success(_) => ()
case Failure(error) => println(s"send failed: ${error.getMessage}")
requestObserver match {
case Some(observer) =>
Try(observer.onNext(request)) match {
case Success(_) => ()
case Failure(error) =>
println(s"send failed: ${error.getMessage}")
reconnect()
}
case None => println("No active stream. Waiting for reconnection.")
}
private def openRequestStream(): Unit = synchronized {
requestObserver = Some(stub.streamUpdates(responseObserver))
}
private def reconnect(): Unit =
if running.get() && reconnecting.compareAndSet(false, true) then {
synchronized { requestObserver = None }
val reconnectThread = new Thread(
() =>
try {
Thread.sleep(1000)
openRequestStream()
state.currentGameId.foreach { gameId =>
println(s"Reconnecting to game $gameId with saved text state.")
stream(gameId)
}
} catch {
case _: InterruptedException => Thread.currentThread.interrupt()
case error: Throwable => println(s"reconnect failed: ${error.getMessage}")
} finally reconnecting.set(false),
"eagle-text-client-reconnect"
)
reconnectThread.setDaemon(true)
reconnectThread.start()
}
private def parseSelected(commandType: String, json: String): Try[SelectedCommandProto] =