blob: b0f56d04cde0640ca17e7ccd72720fe029d34d75 (
plain) (
blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
|
package cask.util
import scala.concurrent.duration.Duration
import scala.concurrent.{Await, ExecutionContext, Promise}
class WsClient(impl: WebsocketBase)
(f: PartialFunction[cask.util.Ws.Event, Unit])
(implicit ec: ExecutionContext, log: Logger)
extends cask.util.BatchActor[Ws.Event]{
def run(items: Seq[Ws.Event]): Unit = items.foreach{
case Ws.Text(s) => impl.send(s)
case Ws.Binary(s) => impl.send(s)
case Ws.Close(_, _) => impl.close()
case Ws.ChannelClosed() => impl.close()
}
def close() = impl.close()
}
object WsClient{
def connect(url: String)
(f: PartialFunction[cask.util.Ws.Event, Unit])
(implicit ec: ExecutionContext, log: Logger): WsClient = {
Await.result(connectAsync(url)(f), Duration.Inf)
}
def connectAsync(url: String)
(f: PartialFunction[cask.util.Ws.Event, Unit])
(implicit ec: ExecutionContext, log: Logger): scala.concurrent.Future[WsClient] = {
object receiveActor extends cask.util.BatchActor[Ws.Event] {
def run(items: Seq[Ws.Event]) = items.collect(f)
}
val p = Promise[WsClient]
val impl = new WebsocketClientImpl(url) {
def onOpen() = {
if (!p.isCompleted) p.success(new WsClient(this)(f))
}
def onMessage(message: String) = {
receiveActor.send(Ws.Text(message))
}
def onMessage(message: Array[Byte]) = {
receiveActor.send(Ws.Binary(message))
}
def onClose(code: Int, reason: String) = {
receiveActor.send(Ws.Close(code, reason))
if (!p.isCompleted) p.success(new WsClient(this)(f))
}
def onError(ex: Exception): Unit = {
receiveActor.send(Ws.Error(ex))
}
}
impl.connect()
p.future
}
}
|