This commit is contained in:
2023-05-04 21:18:04 +08:00
parent 4b65a2459b
commit f971211e33
34 changed files with 90 additions and 87 deletions
+29
View File
@@ -0,0 +1,29 @@
plugins {
kotlin("jvm")
id("com.github.johnrengelman.shadow")
}
group = "cc.maxmc.msm.mastercontrol"
repositories {
mavenCentral()
maven("https://repo.papermc.io/repository/maven-public/")
}
dependencies {
implementation(kotlin("stdlib"))
implementation(project(":common"))
implementation("com.zaxxer:HikariCP:4.0.3")
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-core:1.7.0-Beta")
@Suppress("VulnerableLibrariesLocal")
compileOnly("io.github.waterfallmc:waterfall-api:1.19-R0.1-SNAPSHOT")
}
tasks.shadowJar {
archiveClassifier.set(null as? String?)
relocate("kotlin", "cc.maxmc.msm.lib.kotlin")
}
tasks.build {
dependsOn(tasks.shadowJar)
}
@@ -0,0 +1,33 @@
package cc.maxmc.msm.mastercontrol
import cc.maxmc.msm.api.MultiServerManAPIProvider
import cc.maxmc.msm.mastercontrol.api.APIImpl
import cc.maxmc.msm.mastercontrol.listener.PacketListener
import cc.maxmc.msm.mastercontrol.manager.MatchManager
import cc.maxmc.msm.mastercontrol.manager.ServerManager
import cc.maxmc.msm.mastercontrol.netty.NetManager
import net.md_5.bungee.api.ProxyServer
import net.md_5.bungee.api.plugin.Plugin
class MultiServerMan : Plugin() {
override fun onEnable() {
instance = this
MultiServerManAPIProvider.register(APIImpl)
ProxyServer.getInstance().pluginManager.registerListener(this, PacketListener)
NetManager.startServer()
ServerManager
MatchManager
}
override fun onDisable() {
ServerManager.end()
NetManager.shutdownServer()
}
companion object {
lateinit var instance: MultiServerMan
private set
}
}
@@ -0,0 +1,24 @@
package cc.maxmc.msm.mastercontrol.api
import cc.maxmc.msm.api.MultiServerManAPI
import cc.maxmc.msm.api.misc.MatchInfo
import cc.maxmc.msm.mastercontrol.database.SQLDatabase
import cc.maxmc.msm.mastercontrol.manager.MatchManager
object APIImpl : MultiServerManAPI {
override fun getServer(type: String, players: List<String>): MatchInfo {
return MatchManager.requestMatch(type, players)
}
override fun informEnd(id: Int) {
MatchManager.endMatch(id)
}
override fun getPlayerServer(player: String): MatchInfo {
return SQLDatabase.getPlayerMatch(player)
}
override fun containPlayer(player: String): Boolean {
return getPlayerServer(player).id != -1
}
}
@@ -0,0 +1,88 @@
package cc.maxmc.msm.mastercontrol.database
import cc.maxmc.msm.api.misc.MatchInfo
import cc.maxmc.msm.mastercontrol.manager.MatchManager
import cc.maxmc.msm.mastercontrol.settings.Settings
import com.zaxxer.hikari.HikariConfig
import com.zaxxer.hikari.pool.HikariPool
import java.sql.Timestamp
object SQLDatabase {
private val config = HikariConfig()
private lateinit var pool: HikariPool
fun initDatabase() {
val db = Settings.Database
config.jdbcUrl = "jdbc:mysql://${db.address}:${db.port}"
config.username = db.username
config.password = db.password
config.schema = db.database
pool = HikariPool(config)
}
fun createTable() {
pool.connection.use {
it.createStatement().use { statement ->
statement.execute(
"""
CREATE TABLE `match`
(
`id` int auto_increment,
`start` datetime not null,
`end` datetime null,
`players` text not null ,
PRIMARY KEY (`id`)
)
""".trimIndent()
)
}
}
}
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 MatchInfo()
}
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()): Int =
pool.connection.use {
val prepared = it.prepareStatement(
"""
insert into `match` (type, start, end, players)
values (?, ?, ?, ?);
""".trimIndent()
)
prepared.apply {
setString(1, type)
setTimestamp(2, Timestamp(start))
setTimestamp(3, null)
setString(4, players.joinToString(","))
}
prepared.execute()
prepared.close()
return@use it.createStatement().use { statement ->
val result = statement.executeQuery("select last_insert_id()")
result.next()
result.getInt(1)
}
}
fun endMatch(id: Int, end: Long = System.currentTimeMillis()) = pool.connection.use {
val prepare = it.prepareStatement("update `match` set end = ? where id = ?")
prepare.apply {
setTimestamp(1, Timestamp(end))
setInt(2, id)
}
prepare.execute()
prepare.close()
}
}
@@ -0,0 +1,65 @@
package cc.maxmc.msm.mastercontrol.listener
import cc.maxmc.msm.api.MultiServerManAPIProvider
import cc.maxmc.msm.common.event.ChannelActiveEvent
import cc.maxmc.msm.common.event.ChannelInactiveEvent
import cc.maxmc.msm.common.event.PacketReceiveEvent
import cc.maxmc.msm.common.network.packet.CPacketAPICallback
import cc.maxmc.msm.common.network.packet.CPacketDebug
import cc.maxmc.msm.common.network.packet.PPacketAPICall
import cc.maxmc.msm.common.network.packet.PPacketDebug
import cc.maxmc.msm.common.utils.log
import cc.maxmc.msm.mastercontrol.manager.ChildManager
import net.md_5.bungee.api.plugin.Listener
import net.md_5.bungee.event.EventHandler
object PacketListener : Listener {
val api = MultiServerManAPIProvider.getAPI()
@EventHandler
fun onAPICall(evt: PacketReceiveEvent) {
val packet = evt.packet
if (packet !is PPacketAPICall) return
val callback = when (packet) {
is PPacketAPICall.PPacketCallContainPlayer -> {
CPacketAPICallback.CPacketCallbackContainPlayer(api.containPlayer(packet.player), packet.uid)
}
is PPacketAPICall.PPacketCallGetPlayerServer -> {
CPacketAPICallback.CPacketCallbackGetPlayerServer(api.getPlayerServer(packet.player), packet.uid)
}
is PPacketAPICall.PPacketCallGetServer -> {
CPacketAPICallback.CPacketCallbackGetServer(api.getServer(packet.type, packet.players), packet.uid)
}
is PPacketAPICall.PPacketCallInformEnd -> {
api.informEnd(packet.matchID)
CPacketAPICallback.CPacketCallbackInformEnd(packet.uid)
}
}
evt.channel.writeAndFlush(callback)
}
@EventHandler
fun onPacket(evt: PacketReceiveEvent) {
val packet = evt.packet
if (packet !is PPacketDebug) {
return
}
log("§fDEBUG | §7收到: \"${packet.content}\"")
evt.channel.writeAndFlush(CPacketDebug("(${evt.channel.localAddress()}) - ${packet.content}"))
}
@EventHandler
fun onChannelActive(evt: ChannelActiveEvent) {
log("§a| §7子BC ${evt.channel.remoteAddress()} 成功连接.")
ChildManager.registerChild(evt.channel)
}
@EventHandler
fun onChannelInactive(evt: ChannelInactiveEvent) {
log("§c| §7子BC ${evt.channel.remoteAddress()} 断开连接.")
ChildManager.unregisterChild(evt.channel)
}
}
@@ -0,0 +1,44 @@
package cc.maxmc.msm.mastercontrol.manager
import cc.maxmc.msm.common.network.packet.CPacketGetInfo
import cc.maxmc.msm.common.network.packet.PPacketChildInfo
import cc.maxmc.msm.common.utils.awaitPacket
import cc.maxmc.msm.common.utils.log
import cc.maxmc.msm.common.utils.pluginScope
import cc.maxmc.msm.mastercontrol.misc.ChildBungee
import io.netty.channel.Channel
import kotlinx.coroutines.launch
import java.util.*
import java.util.concurrent.CopyOnWriteArrayList
object ChildManager {
private val children = CopyOnWriteArrayList<ChildBungee>()
fun registerChild(channel: Channel) {
pluginScope.launch {
log("§b| §7正在将 ${channel.remoteAddress()} 注册到集群.")
channel.writeAndFlush(CPacketGetInfo())
val packet = awaitPacket(PPacketChildInfo::class.java) { ch, _ ->
ch == channel
}
val child = ChildBungee(channel, packet.portRange, packet.types)
children.add(child)
ServerManager.initChild(child)
log("§a| §7成功将 ${channel.remoteAddress()} 注册到集群!")
}
}
fun unregisterChild(channel: Channel) {
children.removeIf { it.channel == channel }
}
fun requestChild(type: String): ChildBungee {
val uid = UUID.randomUUID()
children.forEach {
val ports = it.getAvailablePorts()
log("[${uid.toString().substring(0..8)}] ${it.channel.remoteAddress()} has ${ports.size} ports: $ports")
}
return children.filter { it.getAvailablePorts().isNotEmpty() && it.types.contains(type) }
.maxByOrNull { it.getAvailablePorts().size } ?: throw IllegalStateException("当前无可用端口开启新服务器.")
}
}
@@ -0,0 +1,27 @@
package cc.maxmc.msm.mastercontrol.manager
import cc.maxmc.msm.api.misc.MatchInfo
import cc.maxmc.msm.mastercontrol.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) ?: return MatchInfo()
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)
}
}
@@ -0,0 +1,89 @@
package cc.maxmc.msm.mastercontrol.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.mastercontrol.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) {
repeat(2) { _ ->
requestServer(it)
}
}
}
}
}
fun consumeServer(type: String): ServerInfo? {
val servers = getAvailableServers(type)
val info = servers.firstOrNull() ?: return null
info.isAvailable = false
if (servers.size - 1 <= 2) {
pluginScope.launch {
log("§b| §7剩余 $type 服务器不足,正在启动新服务端。")
val server = requestServer(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)
val child = childMap.remove(result)!!
result?.let { child.sendPacket(CPacketEndServer(it)) }
}
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 suspend fun requestServer(type: String): ServerInfo {
val child = ChildManager.requestChild(type)
val list = serverMap[type]!!
val server = child.requestServer(type)
if (server.uid == ServerInfo.NULL_UID)
list.add(server)
childMap[server] = child
return server
}
}
@@ -0,0 +1,41 @@
package cc.maxmc.msm.mastercontrol.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 cc.maxmc.msm.common.utils.log
import io.netty.channel.Channel
import java.net.InetSocketAddress
import java.util.*
class ChildBungee(
val channel: Channel, val ports: Set<Int>, val types: Set<String>, val usedPorts: MutableSet<Int> = HashSet()
) {
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}:$port"
)
val cPacket = CPacketRequestServer(type, server)
sendPacket(cPacket)
val packet = awaitPacket(PPacketServerStarted::class.java) { ch, packet ->
ch == channel && packet.serverInfo == server
}
if (packet.serverInfo.server == null) {
log("§c| §7服务器启动失败,请检查子BC端 §c(${channel.remoteAddress()})§7 日志。")
return ServerInfo()
}
usedPorts += port
return packet.serverInfo
}
fun getAvailablePorts() = ports - usedPorts
}
@@ -0,0 +1,33 @@
package cc.maxmc.msm.mastercontrol.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.mastercontrol.settings.Settings
import io.netty.bootstrap.ServerBootstrap
import io.netty.channel.ChannelFutureListener
import io.netty.channel.nio.NioEventLoopGroup
import io.netty.channel.socket.nio.NioServerSocketChannel
object NetManager {
private val parentGroup = NioEventLoopGroup()
private val childGroup = NioEventLoopGroup()
fun startServer() {
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集群主服务端启动成功. ${
it.channel().localAddress()
}"
)
result.printStackTrace()
})
}
fun shutdownServer() {
parentGroup.shutdownGracefully().sync()
childGroup.shutdownGracefully().sync()
}
}
@@ -0,0 +1,21 @@
package cc.maxmc.msm.mastercontrol.settings
import cc.maxmc.msm.mastercontrol.settings.SettingsReader.config
object Settings {
val serverPort
get() = config.getInt("manage_port", 23333)
object Database {
val address: String
get() = config.getString("database.address", "localhost")
val port: Int
get() = config.getInt("database.port", 12345)
val username: String
get() = config.getString("database.username", "root")
val password: String
get() = config.getString("database.password", "password")
val database: String
get() = config.getString("database.database", "multiserverman")
}
}
@@ -0,0 +1,25 @@
package cc.maxmc.msm.mastercontrol.settings
import cc.maxmc.msm.mastercontrol.MultiServerMan
import net.md_5.bungee.api.ProxyServer
import net.md_5.bungee.api.chat.TextComponent
import net.md_5.bungee.config.Configuration
import net.md_5.bungee.config.ConfigurationProvider
import net.md_5.bungee.config.YamlConfiguration
import kotlin.io.path.*
object SettingsReader {
private val file = MultiServerMan.instance.dataFolder.toPath().resolve("settings.yml")
val config: Configuration
init {
if (!file.exists()) {
MultiServerMan.instance.dataFolder.toPath().createDirectories()
val stream = MultiServerMan.instance.getResourceAsStream("settings.yml")
file.createFile()
stream.copyTo(file.outputStream())
ProxyServer.getInstance().console.sendMessage(TextComponent("§b| §7配置文件不存在,正在创建配置文件。"))
}
config = ConfigurationProvider.getProvider(YamlConfiguration::class.java).load(file.inputStream())
}
}
@@ -0,0 +1,3 @@
name: MasterControl
main: cc.maxmc.msm.mastercontrol.MultiServerMan
author: MistyRain
@@ -0,0 +1,10 @@
# 子BC连接主BC的端口
managePort: 23333
# 数据库配置
database:
address: 127.0.0.1
port: 3306
username: "root"
password: "password"
database: "multiserverman"