Increases timeout period and avoid double websocket connections.
This commit is contained in:
@@ -17,12 +17,14 @@ class Relay(
|
|||||||
var write: Boolean = true
|
var write: Boolean = true
|
||||||
) {
|
) {
|
||||||
private val httpClient = OkHttpClient.Builder()
|
private val httpClient = OkHttpClient.Builder()
|
||||||
.connectTimeout(10, TimeUnit.SECONDS)
|
.connectTimeout(100, TimeUnit.SECONDS)
|
||||||
.readTimeout(10, TimeUnit.SECONDS)
|
.readTimeout(100, TimeUnit.SECONDS)
|
||||||
|
.callTimeout(100, TimeUnit.SECONDS)
|
||||||
.build();
|
.build();
|
||||||
|
|
||||||
private var listeners = setOf<Listener>()
|
private var listeners = setOf<Listener>()
|
||||||
private var socket: WebSocket? = null
|
private var socket: WebSocket? = null
|
||||||
|
private var isReady: Boolean = false
|
||||||
|
|
||||||
var eventDownloadCounter = 0
|
var eventDownloadCounter = 0
|
||||||
var eventUploadCounter = 0
|
var eventUploadCounter = 0
|
||||||
@@ -43,14 +45,19 @@ class Relay(
|
|||||||
return socket != null
|
return socket != null
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Synchronized
|
||||||
fun requestAndWatch() {
|
fun requestAndWatch() {
|
||||||
|
if (socket != null) return
|
||||||
|
|
||||||
try {
|
try {
|
||||||
val request = Request.Builder().url(url.trim()).build()
|
val request = Request.Builder().url(url.trim()).build()
|
||||||
val listener = object : WebSocketListener() {
|
val listener = object : WebSocketListener() {
|
||||||
|
|
||||||
override fun onOpen(webSocket: WebSocket, response: Response) {
|
override fun onOpen(webSocket: WebSocket, response: Response) {
|
||||||
|
isReady = true
|
||||||
ping = response.receivedResponseAtMillis - response.sentRequestAtMillis
|
ping = response.receivedResponseAtMillis - response.sentRequestAtMillis
|
||||||
|
|
||||||
|
// Log.w("Relay", "Relay OnOpen, Loading All subscriptions $url")
|
||||||
// Sends everything.
|
// Sends everything.
|
||||||
Client.allSubscriptions().forEach {
|
Client.allSubscriptions().forEach {
|
||||||
sendFilter(requestId = it)
|
sendFilter(requestId = it)
|
||||||
@@ -65,21 +72,26 @@ class Relay(
|
|||||||
val channel = msg[1].asString
|
val channel = msg[1].asString
|
||||||
when (type) {
|
when (type) {
|
||||||
"EVENT" -> {
|
"EVENT" -> {
|
||||||
|
//Log.w("Relay", "Relay onEVENT $url, $channel")
|
||||||
eventDownloadCounter++
|
eventDownloadCounter++
|
||||||
val event = Event.fromJson(msg[2], Client.lenient)
|
val event = Event.fromJson(msg[2], Client.lenient)
|
||||||
listeners.forEach { it.onEvent(this@Relay, channel, event) }
|
listeners.forEach { it.onEvent(this@Relay, channel, event) }
|
||||||
}
|
}
|
||||||
"EOSE" -> listeners.forEach {
|
"EOSE" -> listeners.forEach {
|
||||||
|
//Log.w("Relay", "Relay onEOSE $url, $channel")
|
||||||
it.onRelayStateChange(this@Relay, Type.EOSE, channel)
|
it.onRelayStateChange(this@Relay, Type.EOSE, channel)
|
||||||
}
|
}
|
||||||
"NOTICE" -> listeners.forEach {
|
"NOTICE" -> listeners.forEach {
|
||||||
|
//Log.w("Relay", "Relay onNotice $url, $channel")
|
||||||
// "channel" being the second string in the string array ...
|
// "channel" being the second string in the string array ...
|
||||||
it.onError(this@Relay, channel, Error("Relay sent notice: $channel"))
|
it.onError(this@Relay, channel, Error("Relay sent notice: $channel"))
|
||||||
}
|
}
|
||||||
"OK" -> listeners.forEach {
|
"OK" -> listeners.forEach {
|
||||||
|
//Log.w("Relay", "Relay onOK $url, $channel")
|
||||||
it.onSendResponse(this@Relay, msg[1].asString, msg[2].asBoolean, msg[3].asString)
|
it.onSendResponse(this@Relay, msg[1].asString, msg[2].asBoolean, msg[3].asString)
|
||||||
}
|
}
|
||||||
else -> listeners.forEach {
|
else -> listeners.forEach {
|
||||||
|
//Log.w("Relay", "Relay something else $url, $channel")
|
||||||
it.onError(
|
it.onError(
|
||||||
this@Relay,
|
this@Relay,
|
||||||
channel,
|
channel,
|
||||||
@@ -105,6 +117,7 @@ class Relay(
|
|||||||
|
|
||||||
override fun onClosed(webSocket: WebSocket, code: Int, reason: String) {
|
override fun onClosed(webSocket: WebSocket, code: Int, reason: String) {
|
||||||
socket = null
|
socket = null
|
||||||
|
isReady = false
|
||||||
closingTime = Date().time / 1000
|
closingTime = Date().time / 1000
|
||||||
listeners.forEach { it.onRelayStateChange(this@Relay, Type.DISCONNECT, null) }
|
listeners.forEach { it.onRelayStateChange(this@Relay, Type.DISCONNECT, null) }
|
||||||
}
|
}
|
||||||
@@ -115,6 +128,7 @@ class Relay(
|
|||||||
socket?.close(1000, "Normal close")
|
socket?.close(1000, "Normal close")
|
||||||
// Failures disconnect the relay.
|
// Failures disconnect the relay.
|
||||||
socket = null
|
socket = null
|
||||||
|
isReady = false
|
||||||
closingTime = Date().time / 1000
|
closingTime = Date().time / 1000
|
||||||
|
|
||||||
Log.w("Relay", "Relay onFailure $url, ${response?.message}")
|
Log.w("Relay", "Relay onFailure $url, ${response?.message}")
|
||||||
@@ -138,6 +152,7 @@ class Relay(
|
|||||||
closingTime = Date().time / 1000
|
closingTime = Date().time / 1000
|
||||||
socket?.close(1000, "Normal close")
|
socket?.close(1000, "Normal close")
|
||||||
socket = null
|
socket = null
|
||||||
|
isReady = false
|
||||||
}
|
}
|
||||||
|
|
||||||
fun sendFilter(requestId: String) {
|
fun sendFilter(requestId: String) {
|
||||||
@@ -148,6 +163,7 @@ class Relay(
|
|||||||
requestAndWatch()
|
requestAndWatch()
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
|
if (isReady) {
|
||||||
val filters = Client.getSubscriptionFilters(requestId)
|
val filters = Client.getSubscriptionFilters(requestId)
|
||||||
if (filters.isNotEmpty()) {
|
if (filters.isNotEmpty()) {
|
||||||
val request =
|
val request =
|
||||||
@@ -158,9 +174,11 @@ class Relay(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fun sendFilterOnlyIfDisconnected() {
|
fun sendFilterOnlyIfDisconnected() {
|
||||||
if (socket == null) {
|
if (socket == null) {
|
||||||
|
println("sendfilter Only if Disconnected ${url} ")
|
||||||
requestAndWatch()
|
requestAndWatch()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user