Archived
0

создана абстракция над NettyServer

This commit is contained in:
2021-04-27 00:10:55 +03:00
parent eff6dcf148
commit 2a1992b6dd
19 changed files with 194 additions and 201 deletions

View File

@@ -1,7 +1,8 @@
apply from: rootDir.toPath().resolve('logic.gradle').toFile()
dependencies {
implementation 'io.netty:netty-all:4.1.22.Final'
implementation libs.netty
implementation libs.reactor
implementation libs.guava
testImplementation libs.lang3

View File

@@ -0,0 +1,20 @@
package mc.protocol;
import io.netty.channel.ChannelHandlerContext;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import mc.protocol.packets.Packet;
@RequiredArgsConstructor
public class ChannelContext<P extends Packet> {
@Getter
private final ChannelHandlerContext ctx;
@Getter
private final P packet;
public void setState(State state) {
ctx.channel().attr(NetworkAttributes.STATE).set(state);
}
}

View File

@@ -0,0 +1,43 @@
package mc.protocol;
import io.netty.bootstrap.ServerBootstrap;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import mc.protocol.di.DaggerProtocolComponent;
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
@RequiredArgsConstructor
public class NettyServer {
private final ServerBootstrap serverBootstrap;
private final Map<Class<? extends Packet>, Sinks.Many<ChannelContext>> observedMap;
public void bind(String host, int port) {
log.info("Network starting: {}:{}", host, port);
try {
serverBootstrap.bind(host, port).sync().channel().closeFuture().sync();
} catch (InterruptedException e) {
if (log.isTraceEnabled()) {
log.trace("{}: {}", e.getClass().getSimpleName(), e.getMessage(), e);
}
}
}
@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() {
ProtocolComponent component = DaggerProtocolComponent.create();
return component.getNettyServer();
}
}

View File

@@ -0,0 +1,21 @@
package mc.protocol;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import lombok.RequiredArgsConstructor;
import mc.protocol.packets.Packet;
import reactor.core.publisher.Sinks;
import java.util.Map;
@SuppressWarnings("rawtypes")
@RequiredArgsConstructor
public class PacketInboundHandler extends SimpleChannelInboundHandler<Packet> {
private final Map<Class<? extends Packet>, Sinks.Many<ChannelContext>> observedMap;
@Override
protected void channelRead0(ChannelHandlerContext ctx, Packet packet) {
observedMap.get(packet.getClass()).tryEmitNext(new ChannelContext<>(ctx, packet));
}
}

View File

@@ -55,6 +55,7 @@ public enum State {
@Getter
private final int id;
@Getter
private final BiMap<Integer, Class<? extends Packet>> serverBoundPackets;
private final BiMap<Integer, Class<? extends Packet>> clientBoundPackets;

View File

@@ -0,0 +1,11 @@
package mc.protocol.di;
import dagger.Component;
import mc.protocol.NettyServer;
@Component(modules = ProtocolModule.class)
@ServerScope
public interface ProtocolComponent {
NettyServer getNettyServer();
}

View File

@@ -0,0 +1,92 @@
package mc.protocol.di;
import com.google.common.collect.ImmutableMap;
import dagger.Module;
import dagger.Provides;
import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.ChannelPipeline;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioServerSocketChannel;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
import mc.protocol.ChannelContext;
import mc.protocol.NettyServer;
import mc.protocol.PacketInboundHandler;
import mc.protocol.State;
import mc.protocol.io.codec.ProtocolDecoder;
import mc.protocol.io.codec.ProtocolEncoder;
import mc.protocol.io.codec.ProtocolSplitter;
import mc.protocol.packets.Packet;
import reactor.core.publisher.Sinks;
import javax.inject.Provider;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.stream.Stream;
@Module
public class ProtocolModule {
@SuppressWarnings("rawtypes")
@Provides
NettyServer provideServer(ServerBootstrap serverBootstrap,
Map<Class<? extends Packet>, Sinks.Many<ChannelContext>> observedMap) {
return new NettyServer(serverBootstrap, observedMap);
}
@Provides
ServerBootstrap provideServerBootstrap(ChannelInitializer<SocketChannel> channelChannelInitializer) {
ServerBootstrap bootstrap = new ServerBootstrap();
bootstrap.group(new NioEventLoopGroup(1), new NioEventLoopGroup())
.channel(NioServerSocketChannel.class)
.childHandler(channelChannelInitializer);
return bootstrap;
}
@Provides
ChannelInitializer<SocketChannel> provideChannelChannelInitializer(
Provider<Map<String, ChannelHandler>> channelHandlerMapProvider) {
return new ChannelInitializer<>() {
@Override
protected void initChannel(SocketChannel socketChannel) {
ChannelPipeline pipeline = socketChannel.pipeline();
channelHandlerMapProvider.get().forEach(pipeline::addLast);
}
};
}
@SuppressWarnings("rawtypes")
@Provides
Map<String, ChannelHandler> provideChannelHandlerMap(
Map<Class<? extends Packet>, Sinks.Many<ChannelContext>> observedMap) {
Map<String, ChannelHandler> map = new LinkedHashMap<>();
map.put("packet_splitter", new ProtocolSplitter());
map.put("logger", new LoggingHandler(LogLevel.DEBUG));
map.put("packet_decoder", new ProtocolDecoder(true));
map.put("packet_encoder", new ProtocolEncoder());
map.put("packet_handler", new PacketInboundHandler(observedMap));
return map;
}
@SuppressWarnings("rawtypes")
@Provides
@ServerScope
Map<Class<? extends Packet>, Sinks.Many<ChannelContext>> provideObservedMap() {
ImmutableMap.Builder<Class<? extends Packet>, Sinks.Many<ChannelContext>> builder = ImmutableMap.builder();
Stream.of(State.values())
.flatMap(state -> state.getServerBoundPackets().values().stream())
.forEach(packetClass -> builder.put(packetClass, Sinks.many().multicast().directBestEffort()));
return builder.build();
}
}

View File

@@ -0,0 +1,10 @@
package mc.protocol.di;
import javax.inject.Scope;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
@Scope
@Retention(RetentionPolicy.RUNTIME)
public @interface ServerScope {
}