Archived
0

простая реализация EventBus

This commit is contained in:
2021-05-05 20:09:49 +03:00
parent 0f1c9bfb1b
commit 205e813fc4
6 changed files with 77 additions and 17 deletions

View File

@@ -5,12 +5,15 @@ 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.ClientSidePacket;
import mc.protocol.utils.EventBus;
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
public class NettyServer { public class NettyServer {
private final ServerBootstrap serverBootstrap; private final ServerBootstrap serverBootstrap;
private final EventBus eventBus;
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);
@@ -24,6 +27,10 @@ public class NettyServer {
} }
} }
public <P extends ClientSidePacket> void listenPacket(State state, Class<P> packetClass, EventBus.EventHandler<P> eventHandler) {
this.eventBus.subscribe(state, packetClass, eventHandler);
}
public static NettyServer createServer() { public static NettyServer createServer() {
ProtocolComponent component = DaggerProtocolComponent.create(); ProtocolComponent component = DaggerProtocolComponent.create();
return component.getNettyServer(); return component.getNettyServer();

View File

@@ -4,19 +4,16 @@ 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.ClientSidePacket; import mc.protocol.packets.ClientSidePacket;
import reactor.core.publisher.Sinks; import mc.protocol.utils.EventBus;
@RequiredArgsConstructor @RequiredArgsConstructor
public class PacketInboundHandler extends SimpleChannelInboundHandler<ClientSidePacket> { public class PacketInboundHandler extends SimpleChannelInboundHandler<ClientSidePacket> {
@SuppressWarnings("rawtypes") private final EventBus eventBus;
@Override @Override
protected void channelRead0(ChannelHandlerContext ctx, ClientSidePacket packet) { protected void channelRead0(ChannelHandlerContext ctx, ClientSidePacket packet) {
Sinks.Many<ChannelContext> packetSinks = ctx.channel().attr(NetworkAttributes.STATE) State state = ctx.channel().attr(NetworkAttributes.STATE).get();
.get().getPacketSinks(packet.getClass()); eventBus.emit(state, new ChannelContext<>(ctx, packet));
if (packetSinks != null) {
packetSinks.tryEmitNext(new ChannelContext<>(ctx, packet));
}
} }
} }

View File

@@ -16,6 +16,8 @@ import mc.protocol.PacketInboundHandler;
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.utils.EventBus;
import mc.protocol.utils.SimpleEventBus;
import javax.annotation.Nonnull; import javax.annotation.Nonnull;
import javax.inject.Provider; import javax.inject.Provider;
@@ -26,8 +28,8 @@ import java.util.Map;
public class ProtocolModule { public class ProtocolModule {
@Provides @Provides
NettyServer provideServer(ServerBootstrap serverBootstrap) { NettyServer provideServer(ServerBootstrap serverBootstrap, EventBus eventBus) {
return new NettyServer(serverBootstrap); return new NettyServer(serverBootstrap, eventBus);
} }
@Provides @Provides
@@ -55,15 +57,21 @@ public class ProtocolModule {
} }
@Provides @Provides
Map<String, ChannelHandler> provideChannelHandlerMap() { Map<String, ChannelHandler> provideChannelHandlerMap(EventBus eventBus) {
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()); map.put("packet_handler", new PacketInboundHandler(eventBus));
return map; return map;
} }
@Provides
@ServerScope
EventBus provideEventBus() {
return new SimpleEventBus();
}
} }

View File

@@ -0,0 +1,25 @@
package mc.protocol.utils;
import mc.protocol.ChannelContext;
import mc.protocol.State;
import mc.protocol.packets.ClientSidePacket;
@SuppressWarnings({ "rawtypes", "unchecked" })
public class SimpleEventBus implements EventBus {
private final Table<State, Class<? extends ClientSidePacket>, EventHandler> table = new Table<>();
@Override
public <P extends ClientSidePacket> void subscribe(State state, Class<P> packetClass, EventHandler<P> eventHandler) {
table.put(state, packetClass, eventHandler);
}
@Override
public <P extends ClientSidePacket> void emit(State state, ChannelContext<P> channelContext) {
EventHandler eventHandler = table.getColumnAndRow(state, channelContext.getPacket().getClass());
if (eventHandler != null) {
eventHandler.handle(channelContext);
}
}
}

View File

@@ -0,0 +1,23 @@
package mc.protocol.utils;
import javax.annotation.Nullable;
import java.util.HashMap;
import java.util.Map;
public class Table<C, R, V> {
private final Map<C, Map<R, V>> map = new HashMap<>();
@Nullable
public V getColumnAndRow(C column, R row) {
if (!map.containsKey(column)) {
return null;
}
return map.get(column).get(row);
}
public void put(C column, R row, V value) {
map.computeIfAbsent(column, c -> new HashMap<>()).put(row, value);
}
}

View File

@@ -49,11 +49,11 @@ public class Main {
NettyServer server = NettyServer.createServer(); NettyServer server = NettyServer.createServer();
PacketHandler packetHandler = serverComponent.getPacketHandler(); PacketHandler packetHandler = serverComponent.getPacketHandler();
State.HANDSHAKING.packetFlux(HandshakePacket.class).subscribe(packetHandler::onHandshake); server.listenPacket(State.HANDSHAKING, HandshakePacket.class, packetHandler::onHandshake);
State.STATUS.packetFlux(PingPacket.class).subscribe(packetHandler::onKeepAlive); server.listenPacket(State.STATUS, PingPacket.class, packetHandler::onKeepAlive);
State.STATUS.packetFlux(StatusServerRequestPacket.class).subscribe(packetHandler::onServerStatus); server.listenPacket(State.STATUS, StatusServerRequestPacket.class, packetHandler::onServerStatus);
State.LOGIN.packetFlux(LoginStartPacket.class).subscribe(packetHandler::onLoginStart); server.listenPacket(State.LOGIN, LoginStartPacket.class, packetHandler::onLoginStart);
State.PLAY.packetFlux(PingPacket.class).subscribe(packetHandler::onKeepAlivePlay); server.listenPacket(State.PLAY, PingPacket.class, packetHandler::onKeepAlivePlay);
server.bind(config.server().host(), config.server().port()); server.bind(config.server().host(), config.server().port());
} }