2023-06-14 09:31:50 +02:00
|
|
|
/*
|
|
|
|
* Copyright 2023 dorkbox, llc
|
|
|
|
*
|
|
|
|
* Licensed under the Apache License, Version 2.0 (the "License");
|
|
|
|
* you may not use this file except in compliance with the License.
|
|
|
|
* You may obtain a copy of the License at
|
|
|
|
*
|
|
|
|
* http://www.apache.org/licenses/LICENSE-2.0
|
|
|
|
*
|
|
|
|
* Unless required by applicable law or agreed to in writing, software
|
|
|
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
|
|
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
|
|
* See the License for the specific language governing permissions and
|
|
|
|
* limitations under the License.
|
|
|
|
*/
|
|
|
|
|
|
|
|
package dorkboxTest.network
|
|
|
|
|
|
|
|
import dorkbox.network.aeron.AeronDriver
|
|
|
|
import dorkbox.network.aeron.endpoint
|
|
|
|
import dorkbox.network.exceptions.ClientTimedOutException
|
2023-06-28 15:09:14 +02:00
|
|
|
import io.aeron.CommonContext
|
2023-06-14 09:31:50 +02:00
|
|
|
import io.aeron.Publication
|
|
|
|
import kotlinx.coroutines.runBlocking
|
|
|
|
import mu.KotlinLogging
|
2023-06-16 14:15:27 +02:00
|
|
|
import org.junit.Assert
|
2023-06-14 09:31:50 +02:00
|
|
|
import org.junit.Test
|
|
|
|
|
|
|
|
class AeronPubSubTest : BaseTest() {
|
|
|
|
@Test
|
|
|
|
fun connectTest() {
|
|
|
|
runBlocking {
|
|
|
|
val log = KotlinLogging.logger("ConnectTest")
|
|
|
|
|
|
|
|
// NOTE: once a config is assigned to a driver, the config cannot be changed
|
|
|
|
val totalCount = 40
|
|
|
|
val port = 3535
|
|
|
|
val serverStreamId = 55555
|
|
|
|
val handshakeTimeoutSec = 10
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
val serverDriver = run {
|
|
|
|
val conf = serverConfig()
|
|
|
|
conf.enableIPv6 = false
|
|
|
|
conf.uniqueAeronDirectory = true
|
|
|
|
|
|
|
|
val driver = AeronDriver(conf, log)
|
|
|
|
driver.start()
|
|
|
|
driver
|
|
|
|
}
|
|
|
|
|
|
|
|
val clientDrivers = mutableListOf<AeronDriver>()
|
|
|
|
val clientPublications = mutableListOf<Pair<AeronDriver, Publication>>()
|
|
|
|
|
|
|
|
|
|
|
|
for (i in 1..totalCount) {
|
|
|
|
val conf = clientConfig()
|
|
|
|
conf.enableIPv6 = false
|
|
|
|
conf.uniqueAeronDirectory = true
|
|
|
|
|
|
|
|
val driver = AeronDriver(conf, log)
|
|
|
|
driver.start()
|
|
|
|
|
|
|
|
clientDrivers.add(driver)
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
2023-06-28 15:09:14 +02:00
|
|
|
val subscriptionUri = AeronDriver.uriHandshake(CommonContext.UDP_MEDIA, true)
|
|
|
|
.endpoint(true, "127.0.0.1", port)
|
2023-06-14 09:31:50 +02:00
|
|
|
val sub = serverDriver.addSubscription(subscriptionUri, serverStreamId, "server")
|
|
|
|
|
|
|
|
var sessionID = 1234567
|
|
|
|
clientDrivers.forEachIndexed { index, clientDriver ->
|
2023-06-28 15:09:14 +02:00
|
|
|
val publicationUri = AeronDriver.uri(CommonContext.UDP_MEDIA, sessionID++, true)
|
|
|
|
.endpoint(true, "127.0.0.1", port)
|
|
|
|
|
|
|
|
// can throw an exception! We catch it in the calling class
|
|
|
|
val publication = clientDriver.addPublication(publicationUri, serverStreamId, "client_$index")
|
|
|
|
|
|
|
|
// can throw an exception! We catch it in the calling class
|
|
|
|
// we actually have to wait for it to connect before we continue
|
|
|
|
clientDriver.waitForConnection(publication, handshakeTimeoutSec, "client_$index") { cause ->
|
2023-06-14 09:31:50 +02:00
|
|
|
ClientTimedOutException("Client publication cannot connect with localhost server", cause)
|
|
|
|
}
|
2023-06-28 15:09:14 +02:00
|
|
|
|
|
|
|
clientPublications.add(Pair(clientDriver, publication))
|
2023-06-14 09:31:50 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
clientPublications.forEachIndexed { index, (clientDriver, pub) ->
|
|
|
|
clientDriver.closeAndDeletePublication(pub, "client_$index")
|
|
|
|
}
|
|
|
|
|
|
|
|
serverDriver.closeAndDeleteSubscription(sub, "server")
|
|
|
|
|
|
|
|
|
2023-06-14 23:35:22 +02:00
|
|
|
clientDrivers.forEach { clientDriver ->
|
|
|
|
clientDriver.close()
|
|
|
|
}
|
|
|
|
|
|
|
|
clientDrivers.forEach { clientDriver ->
|
|
|
|
clientDriver.ensureStopped(10_000, 500)
|
|
|
|
}
|
|
|
|
|
|
|
|
serverDriver.close()
|
|
|
|
serverDriver.ensureStopped(10_000, 500)
|
2023-06-25 17:25:29 +02:00
|
|
|
|
|
|
|
// have to make sure that the aeron driver is CLOSED.
|
|
|
|
Assert.assertTrue("The aeron drivers are not fully closed!", AeronDriver.areAllInstancesClosed())
|
2023-06-14 23:35:22 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2023-06-16 14:15:27 +02:00
|
|
|
@Test()
|
2023-06-14 23:35:22 +02:00
|
|
|
fun connectFailWithBadSessionIdTest() {
|
|
|
|
runBlocking {
|
|
|
|
val log = KotlinLogging.logger("ConnectTest")
|
|
|
|
|
|
|
|
// NOTE: once a config is assigned to a driver, the config cannot be changed
|
|
|
|
val totalCount = 40
|
|
|
|
val port = 3535
|
|
|
|
val serverStreamId = 55555
|
|
|
|
val handshakeTimeoutSec = 10
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
val serverDriver = run {
|
|
|
|
val conf = serverConfig()
|
|
|
|
conf.enableIPv6 = false
|
|
|
|
conf.uniqueAeronDirectory = true
|
|
|
|
|
|
|
|
val driver = AeronDriver(conf, log)
|
|
|
|
driver.start()
|
|
|
|
driver
|
|
|
|
}
|
|
|
|
|
|
|
|
val clientDrivers = mutableListOf<AeronDriver>()
|
|
|
|
val clientPublications = mutableListOf<Pair<AeronDriver, Publication>>()
|
|
|
|
|
|
|
|
|
|
|
|
for (i in 1..totalCount) {
|
|
|
|
val conf = clientConfig()
|
|
|
|
conf.enableIPv6 = false
|
|
|
|
conf.uniqueAeronDirectory = true
|
|
|
|
|
|
|
|
val driver = AeronDriver(conf, log)
|
|
|
|
driver.start()
|
|
|
|
|
|
|
|
clientDrivers.add(driver)
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
2023-06-28 15:09:14 +02:00
|
|
|
val subscriptionUri = AeronDriver.uriHandshake(CommonContext.UDP_MEDIA, true)
|
|
|
|
.endpoint(true, "127.0.0.1", port)
|
2023-06-14 23:35:22 +02:00
|
|
|
val sub = serverDriver.addSubscription(subscriptionUri, serverStreamId, "server")
|
|
|
|
|
2023-06-16 14:15:27 +02:00
|
|
|
try {
|
|
|
|
var sessionID = 1234567
|
|
|
|
clientDrivers.forEachIndexed { index, clientDriver ->
|
2023-06-28 15:09:14 +02:00
|
|
|
val publicationUri = AeronDriver.uri(CommonContext.UDP_MEDIA, sessionID, true)
|
|
|
|
.endpoint(true, "127.0.0.1", port)
|
|
|
|
|
|
|
|
|
|
|
|
// can throw an exception! We catch it in the calling class
|
|
|
|
val publication = clientDriver.addPublication(publicationUri, serverStreamId, "client_$index")
|
|
|
|
|
|
|
|
// can throw an exception! We catch it in the calling class
|
|
|
|
// we actually have to wait for it to connect before we continue
|
|
|
|
clientDriver.waitForConnection(publication, handshakeTimeoutSec, "client_$index") { cause ->
|
2023-06-16 14:15:27 +02:00
|
|
|
ClientTimedOutException("Client publication cannot connect with localhost server", cause)
|
|
|
|
}
|
2023-06-28 15:09:14 +02:00
|
|
|
|
|
|
|
clientPublications.add(Pair(clientDriver, publication))
|
2023-06-14 23:35:22 +02:00
|
|
|
}
|
2023-06-16 14:15:27 +02:00
|
|
|
Assert.fail("TimeoutException should be caught!")
|
|
|
|
} catch (ignore: Exception) {
|
2023-06-14 23:35:22 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
clientPublications.forEachIndexed { index, (clientDriver, pub) ->
|
|
|
|
clientDriver.closeAndDeletePublication(pub, "client_$index")
|
|
|
|
}
|
|
|
|
|
|
|
|
serverDriver.closeAndDeleteSubscription(sub, "server")
|
|
|
|
|
|
|
|
|
2023-06-14 09:31:50 +02:00
|
|
|
clientDrivers.forEach { clientDriver ->
|
|
|
|
clientDriver.close()
|
|
|
|
}
|
|
|
|
|
|
|
|
clientDrivers.forEach { clientDriver ->
|
|
|
|
clientDriver.ensureStopped(10_000, 500)
|
|
|
|
}
|
|
|
|
|
|
|
|
serverDriver.close()
|
|
|
|
serverDriver.ensureStopped(10_000, 500)
|
2023-06-25 17:25:29 +02:00
|
|
|
|
|
|
|
// have to make sure that the aeron driver is CLOSED.
|
|
|
|
Assert.assertTrue("The aeron drivers are not fully closed!", AeronDriver.areAllInstancesClosed())
|
2023-06-14 09:31:50 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|