diff options
author | Li Haoyi <haoyi.sg@gmail.com> | 2019-09-14 22:37:49 +0800 |
---|---|---|
committer | Li Haoyi <haoyi.sg@gmail.com> | 2019-09-14 22:37:49 +0800 |
commit | 9e58e95add96a075d2cb70aa477441261f481ebd (patch) | |
tree | 243fd851913972de6f0110cd254ffd5d6e3565f5 /cask/src/cask/internal/BatchActor.scala | |
parent | dedeed376e1e906ec2eb1574a73f08c24aba47c8 (diff) | |
download | cask-9e58e95add96a075d2cb70aa477441261f481ebd.tar.gz cask-9e58e95add96a075d2cb70aa477441261f481ebd.tar.bz2 cask-9e58e95add96a075d2cb70aa477441261f481ebd.zip |
First pass at providing a convenient API for handling websockets
Diffstat (limited to 'cask/src/cask/internal/BatchActor.scala')
-rw-r--r-- | cask/src/cask/internal/BatchActor.scala | 37 |
1 files changed, 37 insertions, 0 deletions
diff --git a/cask/src/cask/internal/BatchActor.scala b/cask/src/cask/internal/BatchActor.scala new file mode 100644 index 0000000..1566a18 --- /dev/null +++ b/cask/src/cask/internal/BatchActor.scala @@ -0,0 +1,37 @@ +package cask.internal + +import scala.collection.mutable +import scala.concurrent.ExecutionContext + +/** + * A simple asynchrous actor, allowing safe concurrent asynchronous processing + * of queued items. `run` handles items in batches, to allow for batch + * processing optimizations to be used where relevant. + */ +abstract class BatchActor[T]()(implicit ec: ExecutionContext) { + def run(items: Seq[T]): Unit + + private val queue = new mutable.Queue[T]() + private var scheduled = false + def send(t: => T): Unit = synchronized{ + queue.enqueue(t) + if (!scheduled){ + scheduled = true + ec.execute(() => runWithItems()) + } + } + + def runWithItems(): Unit = { + val items = synchronized(queue.dequeueAll(_ => true)) + try run(items) + catch{case e: Throwable => e.printStackTrace()} + synchronized{ + if (queue.nonEmpty) ec.execute(() => runWithItems()) + else{ + assert(scheduled) + scheduled = false + } + } + + } +} |