This commit is contained in:
tony_all
2023-04-16 17:16:38 +08:00
parent 31109849fe
commit fdaa4a57eb
13 changed files with 100 additions and 16 deletions
+1
View File
@@ -20,6 +20,7 @@ dependencies {
}
tasks.shadowJar {
archiveClassifier.set(null as? String?)
relocate("kotlin", "cc.maxmc.msm.lib.kotlin")
}
@@ -7,6 +7,7 @@ import cc.maxmc.msm.parent.manager.MatchManager
import cc.maxmc.msm.parent.manager.ServerManager
import cc.maxmc.msm.parent.netty.NetManager
import net.md_5.bungee.api.ProxyServer
import net.md_5.bungee.api.connection.Server
import net.md_5.bungee.api.plugin.Plugin
class MultiServerMan : Plugin() {
@@ -21,6 +22,7 @@ class MultiServerMan : Plugin() {
}
override fun onDisable() {
ServerManager.end()
NetManager.shutdownServer()
}
@@ -22,6 +22,8 @@ object ChildManager {
}
val child = ChildBungee(channel, packet.portRange, packet.types)
children.add(child)
ServerManager.initChild(child)
log("§a| §7成功将 ${channel.remoteAddress()} 注册到集群!")
}
}
@@ -30,6 +32,10 @@ object ChildManager {
}
fun requestChild(type: String): ChildBungee {
children.forEach {
val ports = it.getAvailablePorts()
log("${ports.size} ports: $ports")
}
return children.filter { it.getAvailablePorts().isNotEmpty() && it.types.contains(type) }
.maxByOrNull { it.getAvailablePorts().size } ?: throw IllegalStateException("当前无可用端口开启新服务器.")
}
@@ -1,14 +1,50 @@
package cc.maxmc.msm.parent.manager
import cc.maxmc.msm.api.misc.ServerInfo
import cc.maxmc.msm.common.network.packet.CPacketEndServer
import cc.maxmc.msm.common.utils.log
import cc.maxmc.msm.common.utils.pluginScope
import cc.maxmc.msm.parent.misc.ChildBungee
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.launch
import java.util.*
import java.util.concurrent.ConcurrentHashMap
object ServerManager {
private val serverMap = ConcurrentHashMap<String, MutableList<ServerInfo>>()
private val childMap = ConcurrentHashMap<ServerInfo, ChildBungee>()
private val serverChannel = Channel<suspend () -> Unit>()
init {
pluginScope.launch {
for (func in serverChannel) {
log("execute one")
func.invoke()
}
}
}
fun end() {
serverChannel.close()
}
suspend fun initChild(child: ChildBungee) {
serverChannel.send {
log("§b| §7正在初始化 ${child.channel.remoteAddress()}")
child.types.forEach {
if (!serverMap.containsKey(it)) {
serverMap[it] = ArrayList()
}
if (serverMap[it]!!.size <= 2) {
pluginScope.launch {
repeat(2) { _ ->
requestServer(it)
}
}
}
}
}
}
fun consumeServer(type: String): ServerInfo {
val servers = getAvailableServers(type)
@@ -17,7 +53,7 @@ object ServerManager {
if (servers.size - 1 <= 2) {
pluginScope.launch {
log("§b| §7剩余 $type 服务器不足,正在启动新服务端。")
val server = requireServer(type)
val server = requestServer(type)
log("§b| §7类型 $type 服务器启动成功: ${server.server}")
}
}
@@ -30,6 +66,8 @@ object ServerManager {
it.find { server -> server.uid == uid }?.let { server -> result = server; true } ?: false
} ?: throw IllegalStateException("Illegal state.")
list.remove(result)
val child = childMap.remove(result)!!
result?.let { child.sendPacket(CPacketEndServer(it)) }
}
fun getServerById(uid: UUID): ServerInfo? {
@@ -40,11 +78,12 @@ object ServerManager {
return serverMap[type]!!.filter { it.isAvailable }
}
private suspend fun requireServer(type: String): ServerInfo {
private suspend fun requestServer(type: String): ServerInfo {
val child = ChildManager.requestChild(type)
val list = serverMap.getOrPut(type) { ArrayList() }
val list = serverMap[type]!!
val server = child.requestServer(type)
list.add(server)
childMap[server] = child
return server
}
}
@@ -12,21 +12,23 @@ import java.util.*
class ChildBungee(
val channel: Channel, val ports: Set<Int>, val types: Set<String>, val usedPorts: MutableSet<Int> = HashSet()
) {
private fun sendPacket(packet: BungeePacket) {
fun sendPacket(packet: BungeePacket) {
channel.writeAndFlush(packet)
}
suspend fun requestServer(type: String): ServerInfo {
val uid = UUID.randomUUID()
val port = getAvailablePorts().first()
val server = ServerInfo(
uid,
"${(channel.remoteAddress() as InetSocketAddress).hostString}:${getAvailablePorts().first()}"
"${(channel.remoteAddress() as InetSocketAddress).hostString}:$port"
)
val cPacket = CPacketRequestServer(type, server)
sendPacket(cPacket)
val packet = awaitPacket(PPacketServerStarted::class.java) { ch, packet ->
ch == channel && packet.serverInfo == server
}
usedPorts += port
return packet.serverInfo
}