This commit is contained in:
2023-04-16 09:31:37 +08:00
parent c534b9508b
commit 31109849fe
30 changed files with 369 additions and 144 deletions
@@ -3,6 +3,8 @@ package cc.maxmc.msm.parent
import cc.maxmc.msm.api.MultiServerManAPIProvider
import cc.maxmc.msm.parent.api.APIImpl
import cc.maxmc.msm.parent.listener.PacketListener
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.plugin.Plugin
@@ -14,6 +16,8 @@ class MultiServerMan : Plugin() {
MultiServerManAPIProvider.register(APIImpl)
ProxyServer.getInstance().pluginManager.registerListener(this, PacketListener)
NetManager.startServer()
ServerManager
MatchManager
}
override fun onDisable() {
@@ -1,9 +0,0 @@
package cc.maxmc.msm.parent
import sun.misc.Signal
import java.lang.management.ManagementFactory
import kotlin.system.exitProcess
fun main() {
}
@@ -1,27 +1,24 @@
package cc.maxmc.msm.parent.api
import cc.maxmc.msm.api.MultiServerManAPI
import cc.maxmc.msm.api.misc.ServerInfo
import com.google.common.net.HostAndPort
import cc.maxmc.msm.api.misc.MatchInfo
import cc.maxmc.msm.parent.database.SQLDatabase
import cc.maxmc.msm.parent.manager.MatchManager
object APIImpl : MultiServerManAPI {
override fun getServer(type: String, players: MutableList<String>): ServerInfo {
TODO("not implemented")
return ServerInfo(HostAndPort.fromString("127.0.0.1:23456"), 1024)
override fun getServer(type: String, players: List<String>): MatchInfo {
return MatchManager.requestMatch(type, players)
}
override fun informEnd(id: Int) {
TODO("not implemented")
return
MatchManager.endMatch(id)
}
override fun getPlayerServer(player: String): ServerInfo {
TODO("not implemented")
return ServerInfo(HostAndPort.fromString("127.0.0.1:34567"), 1025)
override fun getPlayerServer(player: String): MatchInfo {
return SQLDatabase.getPlayerMatch(player)
}
override fun containPlayer(player: String): Boolean {
TODO("not implemented")
return true
return getPlayerServer(player).id != -1
}
}
@@ -1,11 +1,13 @@
package cc.maxmc.msm.parent.database
import cc.maxmc.msm.api.misc.MatchInfo
import cc.maxmc.msm.parent.manager.MatchManager
import cc.maxmc.msm.parent.settings.Settings
import com.zaxxer.hikari.HikariConfig
import com.zaxxer.hikari.pool.HikariPool
import java.sql.Timestamp
class SQLDatabase {
object SQLDatabase {
private val config = HikariConfig()
private lateinit var pool: HikariPool
fun initDatabase() {
@@ -36,21 +38,22 @@ class SQLDatabase {
}
}
fun getPlayerMatch(player: String): Int {
fun getPlayerMatch(player: String): MatchInfo {
pool.connection.use {
it.prepareStatement("select `id` from `match` where find_in_set(?, players) AND end IS NULL")
.use { prepared ->
prepared.setString(0, player)
val rs = prepared.executeQuery()
if (!rs.next()) {
return -1
return MatchInfo()
}
return rs.getInt(1)
val id = rs.getInt(1)
return MatchManager.getMatchById(id) ?: throw IllegalStateException("Match Not Exist in db")
}
}
}
fun recordMatch(type: String, players: List<String>, start: Long = System.currentTimeMillis()) =
fun recordMatch(type: String, players: List<String>, start: Long = System.currentTimeMillis()): Int =
pool.connection.use {
val prepared = it.prepareStatement(
"""
@@ -53,10 +53,12 @@ object PacketListener : Listener {
@EventHandler
fun onChannelActive(evt: ChannelActiveEvent) {
log("§a| §7子BC ${evt.channel.remoteAddress()} 成功连接.")
ChildManager.registerChild(evt.channel)
}
fun onChannelInactive(evt: ChannelInactiveEvent) {
log("§c| §7子BC ${evt.channel.remoteAddress()} 断开连接.")
ChildManager.unregisterChild(evt.channel)
}
}
@@ -20,10 +20,8 @@ object ChildManager {
val packet = awaitPacket(PPacketChildInfo::class.java) { ch, _ ->
ch == channel
}
val child = ChildBungee(channel, packet.portRange)
val child = ChildBungee(channel, packet.portRange, packet.types)
children.add(child)
}
}
@@ -31,8 +29,8 @@ object ChildManager {
children.removeIf { it.channel == channel }
}
fun requestChild(): ChildBungee {
return children.filter { it.getAvailablePorts().isNotEmpty() }.maxByOrNull { it.getAvailablePorts().size }
?: throw IllegalStateException("当前无可用端口开启新服务器.")
fun requestChild(type: String): ChildBungee {
return children.filter { it.getAvailablePorts().isNotEmpty() && it.types.contains(type) }
.maxByOrNull { it.getAvailablePorts().size } ?: throw IllegalStateException("当前无可用端口开启新服务器.")
}
}
@@ -0,0 +1,27 @@
package cc.maxmc.msm.parent.manager
import cc.maxmc.msm.api.misc.MatchInfo
import cc.maxmc.msm.parent.database.SQLDatabase
import java.util.concurrent.ConcurrentHashMap
object MatchManager {
private val matchMap = ConcurrentHashMap<Int, MatchInfo>()
fun requestMatch(type: String, players: List<String>): MatchInfo {
val server = ServerManager.consumeServer(type)
val id = SQLDatabase.recordMatch(type, players)
val match = MatchInfo(id, server)
matchMap[id] = match
return match
}
fun getMatchById(id: Int): MatchInfo? {
return matchMap[id]
}
fun endMatch(id: Int) {
val match = getMatchById(id) ?: throw IllegalStateException("Match does not exist.")
ServerManager.endServer(match.server.uid)
SQLDatabase.endMatch(id)
}
}
@@ -1,17 +1,50 @@
package cc.maxmc.msm.parent.manager
import cc.maxmc.msm.api.misc.ServerInfo
import cc.maxmc.msm.common.utils.log
import cc.maxmc.msm.common.utils.pluginScope
import kotlinx.coroutines.launch
import java.util.*
import java.util.concurrent.ConcurrentHashMap
object ServerManager {
val serverMap = ConcurrentHashMap<String, List<ServerInfo>>()
private val serverMap = ConcurrentHashMap<String, MutableList<ServerInfo>>()
fun getAvailableServers(type: String): List<ServerInfo> {
fun consumeServer(type: String): ServerInfo {
val servers = getAvailableServers(type)
val info = servers.first()
info.isAvailable = false
if (servers.size - 1 <= 2) {
pluginScope.launch {
log("§b| §7剩余 $type 服务器不足,正在启动新服务端。")
val server = requireServer(type)
log("§b| §7类型 $type 服务器启动成功: ${server.server}")
}
}
return info
}
fun endServer(uid: UUID) {
var result: ServerInfo? = null
val list = serverMap.values.find {
it.find { server -> server.uid == uid }?.let { server -> result = server; true } ?: false
} ?: throw IllegalStateException("Illegal state.")
list.remove(result)
}
fun getServerById(uid: UUID): ServerInfo? {
return serverMap.flatMap { it.value }.find { it.uid == uid }
}
private fun getAvailableServers(type: String): List<ServerInfo> {
return serverMap[type]!!.filter { it.isAvailable }
}
private fun requireServer(type: String) {
val child = ChildManager.requestChild()
private suspend fun requireServer(type: String): ServerInfo {
val child = ChildManager.requestChild(type)
val list = serverMap.getOrPut(type) { ArrayList() }
val server = child.requestServer(type)
list.add(server)
return server
}
}
@@ -1,18 +1,33 @@
package cc.maxmc.msm.parent.misc
import cc.maxmc.msm.api.misc.ServerInfo
import cc.maxmc.msm.common.network.BungeePacket
import cc.maxmc.msm.common.network.packet.CPacketRequestServer
import cc.maxmc.msm.common.network.packet.PPacketServerStarted
import cc.maxmc.msm.common.utils.awaitPacket
import io.netty.channel.Channel
import java.net.InetSocketAddress
import java.util.*
class ChildBungee(val channel: Channel, var ports: Set<Int>, var usedPorts: MutableSet<Int> = HashSet()) {
class ChildBungee(
val channel: Channel, val ports: Set<Int>, val types: Set<String>, val usedPorts: MutableSet<Int> = HashSet()
) {
private fun sendPacket(packet: BungeePacket) {
channel.writeAndFlush(packet)
}
suspend fun requestServer(type: String) {
sendPacket(CPacketRequestServer(type))
awaitPacket()
suspend fun requestServer(type: String): ServerInfo {
val uid = UUID.randomUUID()
val server = ServerInfo(
uid,
"${(channel.remoteAddress() as InetSocketAddress).hostString}:${getAvailablePorts().first()}"
)
val cPacket = CPacketRequestServer(type, server)
sendPacket(cPacket)
val packet = awaitPacket(PPacketServerStarted::class.java) { ch, packet ->
ch == channel && packet.serverInfo == server
}
return packet.serverInfo
}
fun getAvailablePorts() = ports - usedPorts
@@ -3,6 +3,8 @@ package cc.maxmc.msm.parent.netty
import cc.maxmc.msm.common.network.netty.NetworkRegistry
import cc.maxmc.msm.common.utils.log
import cc.maxmc.msm.common.utils.pipelineInit
import cc.maxmc.msm.parent.manager.MatchManager
import cc.maxmc.msm.parent.manager.ServerManager
import cc.maxmc.msm.parent.settings.Settings
import io.netty.bootstrap.ServerBootstrap
import io.netty.channel.ChannelFutureListener
@@ -14,11 +16,8 @@ object NetManager {
private val childGroup = NioEventLoopGroup()
fun startServer() {
ServerBootstrap()
.channel(NioServerSocketChannel::class.java)
.group(parentGroup, childGroup)
.childHandler(pipelineInit(NetworkRegistry.PacketDirection.PARENT_BOUND))
.bind(Settings.serverPort)
ServerBootstrap().channel(NioServerSocketChannel::class.java).group(parentGroup, childGroup)
.childHandler(pipelineInit(NetworkRegistry.PacketDirection.PARENT_BOUND)).bind(Settings.serverPort)
.addListener(ChannelFutureListener {
val result = it.cause() ?: return@ChannelFutureListener log(
"§a| §7集群主服务端启动成功. ${
@@ -33,5 +32,4 @@ object NetManager {
parentGroup.shutdownGracefully().sync()
childGroup.shutdownGracefully().sync()
}
}
@@ -4,7 +4,7 @@ import cc.maxmc.msm.parent.settings.SettingsReader.config
object Settings {
val serverPort
get() = config.getInt("server_port", 25566)
get() = config.getInt("manage_port", 23333)
object Database {
val address: String
+4 -1
View File
@@ -1,4 +1,7 @@
serverPort: 12345
# 子BC连接主BC的端口
managePort: 23333
# 数据库配置
database:
address: 127.0.0.1
port: 3306