I utilized the TimeoutScheduler introduced at Scala Futures - built in timeout?.
However, now my program does not terminate as before without TimeoutScheduler.
I have two Futures: res1 and res2. Both with a timeout of 15 seconds. In the end I sequence both Futures in order to shutdown the HTTP executor properly in the onComplete callback. Without using withTimeout the program terminates right after http.shutdown. But with using withTimeout is doesn't. Why? There must be some further futures...
import java.net.URI
import scala.util.{ Try, Failure, Success }
import dispatch._
import org.json4s._
import org.json4s.native.JsonMethods._
import com.typesafe.scalalogging.slf4j.Logging
object Main extends App with Logging {
import scala.concurrent.ExecutionContext.Implicits.global
val http = Http
val github = host("api.github.com").secure
import timeout._
import scala.concurrent.duration._
import scala.language.postfixOps
import scala.collection.JavaConverters._
val req: dispatch.Req = github / "users" / "defunkt"
val res1 = http(req > dispatch.as.Response(_.getHeaders().get("Server").asScala)) withTimeout (15 seconds) recover { case x =>
logger.debug("Error1: " + x.toString)
Nil
}
val res2 = http(req > dispatch.as.Response(_.getHeaders().get("Vary").asScala)) withTimeout (15 seconds) recover { case x =>
logger.debug("Error2: " + x.toString)
Nil
}
Future.sequence(res1 :: res2 :: Nil) onComplete { case _ =>
http.shutdown() // without using `withTimeout` the program terminated after `http.shutdow`
TimeoutScheduler.timer.stop() // thanks to @cmbaxter
}
}
object timeout {
import java.util.concurrent.TimeUnit
import scala.concurrent.Promise
import scala.concurrent.duration.Duration
import org.jboss.netty.util.Timeout
import org.jboss.netty.util.TimerTask
import org.jboss.netty.util.HashedWheelTimer
import org.jboss.netty.handler.timeout.TimeoutException
// cf. https://stackoverflow.com/questions/16304471/scala-futures-built-in-timeout
object TimeoutScheduler {
val timer = new HashedWheelTimer(10, TimeUnit.MILLISECONDS)
def scheduleTimeout(promise: Promise[_], after: Duration) = {
timer.newTimeout(new TimerTask{
def run(timeout: Timeout){
promise.failure(new TimeoutException("Operation timed out after " + after.toMillis + " millis"))
}
}, after.toNanos, TimeUnit.NANOSECONDS)
}
}
implicit class FutureWithTimeout[T](f: Future[T]) {
import scala.concurrent.ExecutionContext
def withTimeout(after: Duration)(implicit ec: ExecutionContext) = {
val prom = Promise[T]()
val timeout = TimeoutScheduler.scheduleTimeout(prom, after)
val combinedFut = Future.firstCompletedOf(List(f, prom.future))
f onComplete { case result => timeout.cancel() }
combinedFut
}
}
}
Any suggestions are welcome, Best, /nm