Add disconnect lifecycle handlers
This commit is contained in:
parent
7bebbf7bd7
commit
19b349143d
@ -46,6 +46,7 @@ class KiloClientConnection<S>(
|
|||||||
suspend fun run(onConnectedStateChanged: ((Boolean) -> Unit)? = null) {
|
suspend fun run(onConnectedStateChanged: ((Boolean) -> Unit)? = null) {
|
||||||
coroutineScope {
|
coroutineScope {
|
||||||
var job: Job? = null
|
var job: Job? = null
|
||||||
|
var connectedScope: KiloScope<S>? = null
|
||||||
try {
|
try {
|
||||||
// in parallel: keys and connection
|
// in parallel: keys and connection
|
||||||
val deferredKeyPair = async { SafeKeyExchange() }
|
val deferredKeyPair = async { SafeKeyExchange() }
|
||||||
@ -90,7 +91,8 @@ class KiloClientConnection<S>(
|
|||||||
kiloRemoteInterface.complete(
|
kiloRemoteInterface.complete(
|
||||||
KiloRemoteInterface(deferredParams, clientInterface)
|
KiloRemoteInterface(deferredParams, clientInterface)
|
||||||
)
|
)
|
||||||
clientInterface.onConnectHandlers.invokeAll(params.scope)
|
connectedScope = params.scope
|
||||||
|
clientInterface.onConnectHandlers.invokeAll(connectedScope)
|
||||||
onConnectedStateChanged?.invoke(true)
|
onConnectedStateChanged?.invoke(true)
|
||||||
job.join()
|
job.join()
|
||||||
|
|
||||||
@ -99,6 +101,7 @@ class KiloClientConnection<S>(
|
|||||||
} catch (x: RemoteInterface.ClosedException) {
|
} catch (x: RemoteInterface.ClosedException) {
|
||||||
debug { "connection closed/refused by remote" }
|
debug { "connection closed/refused by remote" }
|
||||||
} finally {
|
} finally {
|
||||||
|
connectedScope?.let { clientInterface.onDisconnectHandlers.invokeAll(it) }
|
||||||
onConnectedStateChanged?.invoke(false)
|
onConnectedStateChanged?.invoke(false)
|
||||||
job?.cancel()
|
job?.cancel()
|
||||||
device.apply { runCatching { close() } }
|
device.apply { runCatching { close() } }
|
||||||
|
|||||||
@ -19,16 +19,17 @@ typealias KiloHandler<S> = KiloScope<S>.()->Unit
|
|||||||
*
|
*
|
||||||
* - It registers common exceptions from [RemoteInterface] and kotlin/java `IllegalArgumentException` and
|
* - It registers common exceptions from [RemoteInterface] and kotlin/java `IllegalArgumentException` and
|
||||||
* `IllegalStateException`
|
* `IllegalStateException`
|
||||||
* - It provides [onConnected] handler
|
* - It provides [onConnected] and [onDisconnected] handlers
|
||||||
*
|
*
|
||||||
* See [KiloServer] for usage sample.
|
* See [KiloServer] for usage sample.
|
||||||
*/
|
*/
|
||||||
open class KiloInterface<S> : LocalInterface<KiloScope<S>>() {
|
open class KiloInterface<S> : LocalInterface<KiloScope<S>>() {
|
||||||
|
|
||||||
internal val onConnectHandlers = mutableListOf<KiloHandler<S>>()
|
internal val onConnectHandlers = mutableListOf<KiloHandler<S>>()
|
||||||
|
internal val onDisconnectHandlers = mutableListOf<KiloHandler<S>>()
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Registers handler [f] for [onConnected] event, to the head or the end of handler list.
|
* Registers handler [f] for [onConnected] event, to the head or the end of a handler list.
|
||||||
*
|
*
|
||||||
* @param addFirst if true, [f] will be added to the beginning of the list of handlers
|
* @param addFirst if true, [f] will be added to the beginning of the list of handlers
|
||||||
*/
|
*/
|
||||||
@ -36,6 +37,18 @@ open class KiloInterface<S> : LocalInterface<KiloScope<S>>() {
|
|||||||
if( addFirst ) onConnectHandlers.add(0, f) else onConnectHandlers += f
|
if( addFirst ) onConnectHandlers.add(0, f) else onConnectHandlers += f
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Registers handler [f] for [onDisconnected] event, to the head or the end of a handler list.
|
||||||
|
*
|
||||||
|
* It is called with the same connection [KiloScope] as [onConnected], after an established
|
||||||
|
* connection is closed.
|
||||||
|
*
|
||||||
|
* @param addFirst if true, [f] will be added to the beginning of the list of handlers
|
||||||
|
*/
|
||||||
|
fun onDisconnected(addFirst: Boolean = false, f: KiloScope<S>.()->Unit) {
|
||||||
|
if( addFirst ) onDisconnectHandlers.add(0, f) else onDisconnectHandlers += f
|
||||||
|
}
|
||||||
|
|
||||||
init {
|
init {
|
||||||
registerError { RemoteInterface.UnknownCommand(it) }
|
registerError { RemoteInterface.UnknownCommand(it) }
|
||||||
registerError { RemoteInterface.InternalError(it) }
|
registerError { RemoteInterface.InternalError(it) }
|
||||||
@ -47,4 +60,3 @@ open class KiloInterface<S> : LocalInterface<KiloScope<S>>() {
|
|||||||
registerError { IllegalArgumentException(it) }
|
registerError { IllegalArgumentException(it) }
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@ -53,6 +53,7 @@ class KiloServerConnection<S>(
|
|||||||
suspend fun run() {
|
suspend fun run() {
|
||||||
val deferredParams = CompletableDeferred<KiloParams<S>>()
|
val deferredParams = CompletableDeferred<KiloParams<S>>()
|
||||||
val deferredTransport = CompletableDeferred<Transport<*>>()
|
val deferredTransport = CompletableDeferred<Transport<*>>()
|
||||||
|
var connectedScope: KiloScope<S>? = null
|
||||||
|
|
||||||
val l0Interface = KiloL0Interface(clientInterface, deferredParams).apply {
|
val l0Interface = KiloL0Interface(clientInterface, deferredParams).apply {
|
||||||
var params: KiloParams<S>? = null
|
var params: KiloParams<S>? = null
|
||||||
@ -85,7 +86,8 @@ class KiloServerConnection<S>(
|
|||||||
kiloRemoteInterface.complete(
|
kiloRemoteInterface.complete(
|
||||||
KiloRemoteInterface(deferredParams, clientInterface)
|
KiloRemoteInterface(deferredParams, clientInterface)
|
||||||
)
|
)
|
||||||
clientInterface.onConnectHandlers.invokeAll(p.scope)
|
connectedScope = p.scope
|
||||||
|
clientInterface.onConnectHandlers.invokeAll(connectedScope)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -93,8 +95,12 @@ class KiloServerConnection<S>(
|
|||||||
deferredTransport.complete(transport)
|
deferredTransport.complete(transport)
|
||||||
kiloRemoteInterface.complete(KiloRemoteInterface(deferredParams,clientInterface))
|
kiloRemoteInterface.complete(KiloRemoteInterface(deferredParams,clientInterface))
|
||||||
debug { "starting the transport"}
|
debug { "starting the transport"}
|
||||||
transport.run()
|
try {
|
||||||
debug { "server transport finished" }
|
transport.run()
|
||||||
|
debug { "server transport finished" }
|
||||||
|
} finally {
|
||||||
|
connectedScope?.let { clientInterface.onDisconnectHandlers.invokeAll(it) }
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
companion object {
|
companion object {
|
||||||
|
|||||||
@ -170,6 +170,41 @@ class TransportTest {
|
|||||||
d2.close()
|
d2.close()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun testDisconnectedHandlers() = runTest {
|
||||||
|
initCrypto()
|
||||||
|
|
||||||
|
val cmdPing by command<String, String>()
|
||||||
|
val (d1, d2) = createTestDevice()
|
||||||
|
val serverDisconnected = CompletableDeferred<String>()
|
||||||
|
val clientDisconnected = CompletableDeferred<String>()
|
||||||
|
|
||||||
|
val serverInterface = KiloInterface<String>().apply {
|
||||||
|
onDisconnected {
|
||||||
|
serverDisconnected.complete(session)
|
||||||
|
}
|
||||||
|
on(cmdPing) {
|
||||||
|
"pong! [$it]"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
launch { KiloServerConnection(serverInterface, d1, "server session").run() }
|
||||||
|
|
||||||
|
val clientInterface = KiloInterface<String>().apply {
|
||||||
|
onDisconnected {
|
||||||
|
clientDisconnected.complete(session)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
val client = KiloClientConnection(clientInterface, d2, "client session")
|
||||||
|
launch { client.run() }
|
||||||
|
|
||||||
|
assertEquals("pong! [hello]", client.call(cmdPing, "hello"))
|
||||||
|
d1.close()
|
||||||
|
d2.close()
|
||||||
|
|
||||||
|
assertEquals("server session", withTimeout(1000) { serverDisconnected.await() })
|
||||||
|
assertEquals("client session", withTimeout(1000) { clientDisconnected.await() })
|
||||||
|
}
|
||||||
|
|
||||||
class TestException(text: String) : Exception(text)
|
class TestException(text: String) : Exception(text)
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user