пересмотр событийной модели
This commit is contained in:
@@ -3,10 +3,10 @@ package mc.protocol;
|
|||||||
import io.netty.channel.ChannelHandlerContext;
|
import io.netty.channel.ChannelHandlerContext;
|
||||||
import lombok.Getter;
|
import lombok.Getter;
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import mc.protocol.packets.Packet;
|
import mc.protocol.packets.ClientSidePacket;
|
||||||
|
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class ChannelContext<P extends Packet> {
|
public class ChannelContext<P extends ClientSidePacket> {
|
||||||
|
|
||||||
@Getter
|
@Getter
|
||||||
private final ChannelHandlerContext ctx;
|
private final ChannelHandlerContext ctx;
|
||||||
|
|||||||
@@ -5,19 +5,12 @@ import lombok.RequiredArgsConstructor;
|
|||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import mc.protocol.di.DaggerProtocolComponent;
|
import mc.protocol.di.DaggerProtocolComponent;
|
||||||
import mc.protocol.di.ProtocolComponent;
|
import mc.protocol.di.ProtocolComponent;
|
||||||
import mc.protocol.packets.Packet;
|
|
||||||
import reactor.core.publisher.Flux;
|
|
||||||
import reactor.core.publisher.Sinks;
|
|
||||||
|
|
||||||
import java.util.Map;
|
|
||||||
|
|
||||||
@SuppressWarnings("rawtypes")
|
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class NettyServer {
|
public class NettyServer {
|
||||||
|
|
||||||
private final ServerBootstrap serverBootstrap;
|
private final ServerBootstrap serverBootstrap;
|
||||||
private final Map<Class<? extends Packet>, Sinks.Many<ChannelContext>> observedMap;
|
|
||||||
|
|
||||||
public void bind(String host, int port) {
|
public void bind(String host, int port) {
|
||||||
log.info("Network starting: {}:{}", host, port);
|
log.info("Network starting: {}:{}", host, port);
|
||||||
@@ -31,11 +24,6 @@ public class NettyServer {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@SuppressWarnings("unchecked")
|
|
||||||
public <P extends Packet> Flux<ChannelContext<P>> packetFlux(Class<P> packetClass) {
|
|
||||||
return observedMap.get(packetClass).asFlux().map(ChannelContext.class::cast);
|
|
||||||
}
|
|
||||||
|
|
||||||
public static NettyServer createServer() {
|
public static NettyServer createServer() {
|
||||||
ProtocolComponent component = DaggerProtocolComponent.create();
|
ProtocolComponent component = DaggerProtocolComponent.create();
|
||||||
return component.getNettyServer();
|
return component.getNettyServer();
|
||||||
|
|||||||
@@ -3,21 +3,20 @@ package mc.protocol;
|
|||||||
import io.netty.channel.ChannelHandlerContext;
|
import io.netty.channel.ChannelHandlerContext;
|
||||||
import io.netty.channel.SimpleChannelInboundHandler;
|
import io.netty.channel.SimpleChannelInboundHandler;
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import mc.protocol.packets.Packet;
|
import mc.protocol.packets.ClientSidePacket;
|
||||||
import reactor.core.publisher.Sinks;
|
import reactor.core.publisher.Sinks;
|
||||||
|
|
||||||
import java.util.Map;
|
|
||||||
|
|
||||||
@SuppressWarnings("rawtypes")
|
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class PacketInboundHandler extends SimpleChannelInboundHandler<Packet> {
|
public class PacketInboundHandler extends SimpleChannelInboundHandler<ClientSidePacket> {
|
||||||
|
|
||||||
private final Map<Class<? extends Packet>, Sinks.Many<ChannelContext>> observedMap;
|
|
||||||
|
|
||||||
|
@SuppressWarnings("rawtypes")
|
||||||
@Override
|
@Override
|
||||||
protected void channelRead0(ChannelHandlerContext ctx, Packet packet) {
|
protected void channelRead0(ChannelHandlerContext ctx, ClientSidePacket packet) {
|
||||||
if (observedMap.containsKey(packet.getClass())) {
|
Sinks.Many<ChannelContext> packetSinks = ctx.channel().attr(NetworkAttributes.STATE)
|
||||||
observedMap.get(packet.getClass()).tryEmitNext(new ChannelContext<>(ctx, packet));
|
.get().getPacketSinks(packet.getClass());
|
||||||
|
|
||||||
|
if (packetSinks != null) {
|
||||||
|
packetSinks.tryEmitNext(new ChannelContext<>(ctx, packet));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,9 +8,12 @@ import mc.protocol.packets.PingPacket;
|
|||||||
import mc.protocol.packets.ServerSidePacket;
|
import mc.protocol.packets.ServerSidePacket;
|
||||||
import mc.protocol.packets.client.*;
|
import mc.protocol.packets.client.*;
|
||||||
import mc.protocol.packets.server.*;
|
import mc.protocol.packets.server.*;
|
||||||
|
import reactor.core.publisher.Flux;
|
||||||
|
import reactor.core.publisher.Sinks;
|
||||||
|
|
||||||
import javax.annotation.Nullable;
|
import javax.annotation.Nullable;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
|
import java.util.HashMap;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
@@ -75,10 +78,12 @@ public enum State {
|
|||||||
@Getter
|
@Getter
|
||||||
private final int id;
|
private final int id;
|
||||||
|
|
||||||
@Getter
|
|
||||||
private final Map<Integer, Class<? extends ClientSidePacket>> clientSidePackets;
|
private final Map<Integer, Class<? extends ClientSidePacket>> clientSidePackets;
|
||||||
private final Map<Class<? extends ServerSidePacket>, Integer> serverSidePackets;
|
private final Map<Class<? extends ServerSidePacket>, Integer> serverSidePackets;
|
||||||
|
|
||||||
|
@SuppressWarnings("rawtypes")
|
||||||
|
private final Map<Class<? extends ClientSidePacket>, Sinks.Many<ChannelContext>> observedMap = new HashMap<>();
|
||||||
|
|
||||||
State(int id, Map<Integer, Class<? extends ClientSidePacket>> clientSidePackets) {
|
State(int id, Map<Integer, Class<? extends ClientSidePacket>> clientSidePackets) {
|
||||||
this.id = id;
|
this.id = id;
|
||||||
this.clientSidePackets = clientSidePackets;
|
this.clientSidePackets = clientSidePackets;
|
||||||
@@ -94,4 +99,16 @@ public enum State {
|
|||||||
public Integer getServerSidePacketId(Class<? extends Packet> clazz) {
|
public Integer getServerSidePacketId(Class<? extends Packet> clazz) {
|
||||||
return serverSidePackets == null ? null : serverSidePackets.get(clazz);
|
return serverSidePackets == null ? null : serverSidePackets.get(clazz);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@SuppressWarnings("rawtypes")
|
||||||
|
public <P extends ClientSidePacket> Sinks.Many<ChannelContext> getPacketSinks(Class<P> packetClass) {
|
||||||
|
return observedMap.get(packetClass);
|
||||||
|
}
|
||||||
|
|
||||||
|
@SuppressWarnings("unchecked")
|
||||||
|
public <P extends ClientSidePacket> Flux<ChannelContext<P>> packetFlux(Class<P> packetClass) {
|
||||||
|
return observedMap.computeIfAbsent(packetClass, aClass -> Sinks.many().multicast().directBestEffort())
|
||||||
|
.asFlux().map(ChannelContext.class::cast);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -11,31 +11,23 @@ import io.netty.channel.socket.SocketChannel;
|
|||||||
import io.netty.channel.socket.nio.NioServerSocketChannel;
|
import io.netty.channel.socket.nio.NioServerSocketChannel;
|
||||||
import io.netty.handler.logging.LogLevel;
|
import io.netty.handler.logging.LogLevel;
|
||||||
import io.netty.handler.logging.LoggingHandler;
|
import io.netty.handler.logging.LoggingHandler;
|
||||||
import mc.protocol.ChannelContext;
|
|
||||||
import mc.protocol.NettyServer;
|
import mc.protocol.NettyServer;
|
||||||
import mc.protocol.PacketInboundHandler;
|
import mc.protocol.PacketInboundHandler;
|
||||||
import mc.protocol.State;
|
|
||||||
import mc.protocol.io.codec.ProtocolDecoder;
|
import mc.protocol.io.codec.ProtocolDecoder;
|
||||||
import mc.protocol.io.codec.ProtocolEncoder;
|
import mc.protocol.io.codec.ProtocolEncoder;
|
||||||
import mc.protocol.io.codec.ProtocolSplitter;
|
import mc.protocol.io.codec.ProtocolSplitter;
|
||||||
import mc.protocol.packets.Packet;
|
|
||||||
import reactor.core.publisher.Sinks;
|
|
||||||
|
|
||||||
import javax.annotation.Nonnull;
|
import javax.annotation.Nonnull;
|
||||||
import javax.inject.Provider;
|
import javax.inject.Provider;
|
||||||
import java.util.LinkedHashMap;
|
import java.util.LinkedHashMap;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.stream.Collectors;
|
|
||||||
import java.util.stream.Stream;
|
|
||||||
|
|
||||||
@Module
|
@Module
|
||||||
public class ProtocolModule {
|
public class ProtocolModule {
|
||||||
|
|
||||||
@SuppressWarnings("rawtypes")
|
|
||||||
@Provides
|
@Provides
|
||||||
NettyServer provideServer(ServerBootstrap serverBootstrap,
|
NettyServer provideServer(ServerBootstrap serverBootstrap) {
|
||||||
Map<Class<? extends Packet>, Sinks.Many<ChannelContext>> observedMap) {
|
return new NettyServer(serverBootstrap);
|
||||||
return new NettyServer(serverBootstrap, observedMap);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Provides
|
@Provides
|
||||||
@@ -62,28 +54,16 @@ public class ProtocolModule {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
@SuppressWarnings("rawtypes")
|
|
||||||
@Provides
|
@Provides
|
||||||
Map<String, ChannelHandler> provideChannelHandlerMap(
|
Map<String, ChannelHandler> provideChannelHandlerMap() {
|
||||||
Map<Class<? extends Packet>, Sinks.Many<ChannelContext>> observedMap) {
|
|
||||||
|
|
||||||
Map<String, ChannelHandler> map = new LinkedHashMap<>();
|
Map<String, ChannelHandler> map = new LinkedHashMap<>();
|
||||||
|
|
||||||
map.put("packet_splitter", new ProtocolSplitter());
|
map.put("packet_splitter", new ProtocolSplitter());
|
||||||
map.put("logger", new LoggingHandler(LogLevel.DEBUG));
|
map.put("logger", new LoggingHandler(LogLevel.DEBUG));
|
||||||
map.put("packet_decoder", new ProtocolDecoder(true));
|
map.put("packet_decoder", new ProtocolDecoder(true));
|
||||||
map.put("packet_encoder", new ProtocolEncoder());
|
map.put("packet_encoder", new ProtocolEncoder());
|
||||||
map.put("packet_handler", new PacketInboundHandler(observedMap));
|
map.put("packet_handler", new PacketInboundHandler());
|
||||||
|
|
||||||
return map;
|
return map;
|
||||||
}
|
}
|
||||||
|
|
||||||
@SuppressWarnings("rawtypes")
|
|
||||||
@Provides
|
|
||||||
@ServerScope
|
|
||||||
Map<Class<? extends Packet>, Sinks.Many<ChannelContext>> provideObservedMap() {
|
|
||||||
return Stream.of(State.values())
|
|
||||||
.flatMap(state -> state.getClientSidePackets().values().stream())
|
|
||||||
.collect(Collectors.toMap(packetClass -> packetClass, v -> Sinks.many().multicast().directBestEffort()));
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import joptsimple.OptionSet;
|
|||||||
import joptsimple.util.PathConverter;
|
import joptsimple.util.PathConverter;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import mc.protocol.NettyServer;
|
import mc.protocol.NettyServer;
|
||||||
|
import mc.protocol.State;
|
||||||
import mc.protocol.packets.PingPacket;
|
import mc.protocol.packets.PingPacket;
|
||||||
import mc.protocol.packets.client.HandshakePacket;
|
import mc.protocol.packets.client.HandshakePacket;
|
||||||
import mc.protocol.packets.client.LoginStartPacket;
|
import mc.protocol.packets.client.LoginStartPacket;
|
||||||
@@ -48,10 +49,10 @@ public class Main {
|
|||||||
NettyServer server = NettyServer.createServer();
|
NettyServer server = NettyServer.createServer();
|
||||||
PacketHandler packetHandler = serverComponent.getPacketHandler();
|
PacketHandler packetHandler = serverComponent.getPacketHandler();
|
||||||
|
|
||||||
server.packetFlux(HandshakePacket.class).subscribe(packetHandler::onHandshake);
|
State.HANDSHAKING.packetFlux(HandshakePacket.class).subscribe(packetHandler::onHandshake);
|
||||||
server.packetFlux(PingPacket.class).subscribe(packetHandler::onKeepAlive);
|
State.STATUS.packetFlux(PingPacket.class).subscribe(packetHandler::onKeepAlive);
|
||||||
server.packetFlux(StatusServerRequestPacket.class).subscribe(packetHandler::onServerStatus);
|
State.STATUS.packetFlux(StatusServerRequestPacket.class).subscribe(packetHandler::onServerStatus);
|
||||||
server.packetFlux(LoginStartPacket.class).subscribe(packetHandler::onLoginStart);
|
State.LOGIN.packetFlux(LoginStartPacket.class).subscribe(packetHandler::onLoginStart);
|
||||||
|
|
||||||
server.bind(config.server().host(), config.server().port());
|
server.bind(config.server().host(), config.server().port());
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user