Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
import org.jetlinks.community.network.security.VertxKeyCertTrustOptions;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

import javax.annotation.Nonnull;
Expand Down Expand Up @@ -123,27 +124,31 @@ private Mono<Network> initServer(VertxHttpServer server, HttpServerConfig config
int numberOfInstance = Math.max(1, config.getInstance());
List<HttpServer> instances = new ArrayList<>(numberOfInstance);
return convert(config)
.map(options -> {
.flatMap(options -> {
//利用多线程处理请求
for (int i = 0; i < numberOfInstance; i++) {
instances.add(createHttpServer(options));
}
server.setBindAddress(new InetSocketAddress(config.getHost(), config.getPort()));
server.setLastError(null);
server.setHttpServers(instances);
for (HttpServer httpServer : instances) {
vertx.nettyEventLoopGroup()
.execute(()->{
httpServer.listen(result -> {
if (result.succeeded()) {
log.debug("startup http server on [{}]", server.getBindAddress());
} else {
server.setLastError(result.cause().getMessage());
log.warn("startup http server on [{}] failed", server.getBindAddress(), result.cause());
}
});
});
}
return server;
return Flux
.fromIterable(instances)
.flatMap(httpServer -> Mono
.fromCompletionStage(httpServer.listen().toCompletionStage())
.doOnNext(ignore -> log.debug("startup http server on [{}]", server.getBindAddress())))
.doOnCancel(server::shutdown)
.then(Mono.fromSupplier(() -> {
server.startupComplete();
return server;
}))
.onErrorResume(error -> {
server.setLastError(error.getMessage());
log.warn("startup http server on [{}] failed", server.getBindAddress(), error);
return server
.shutdownAsync()
.then(Mono.error(error));
});
});
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,11 +35,14 @@
import reactor.core.Disposables;
import reactor.core.publisher.Flux;
import reactor.core.publisher.FluxSink;
import reactor.core.publisher.Mono;

import java.net.InetSocketAddress;
import java.util.Collection;
import java.util.Collections;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLongFieldUpdater;
import java.util.stream.Collectors;
import java.util.stream.Stream;
Expand All @@ -58,7 +61,9 @@ public class VertxHttpServer implements HttpServer {

private static final Map<HttpMethod, SeparatedCharSequence> HTTP_PREFIX_CACHE = new ConcurrentHashMap<>();

private Collection<io.vertx.core.http.HttpServer> httpServers;
private volatile Collection<io.vertx.core.http.HttpServer> httpServers;

private final AtomicBoolean started = new AtomicBoolean();

private HttpServerConfig config;

Expand Down Expand Up @@ -91,7 +96,8 @@ public InetSocketAddress getBindAddress() {
}

public void setHttpServers(Collection<io.vertx.core.http.HttpServer> httpServers) {
if (isAlive()) {
started.set(false);
if (this.httpServers != null && !this.httpServers.isEmpty()) {
shutdown();
}
this.httpServers = httpServers;
Expand Down Expand Up @@ -216,6 +222,29 @@ public void setHttpServers(Collection<io.vertx.core.http.HttpServer> httpServers
}
}

void startupComplete() {
started.set(true);
}

Mono<Void> shutdownAsync() {
return Flux
.fromIterable(clearServers())
.flatMap(httpServer -> Mono
.fromCompletionStage(httpServer.close().toCompletionStage())
.onErrorResume(error -> {
log.warn("close http server error", error);
return Mono.empty();
}))
.then();
}

private synchronized Collection<io.vertx.core.http.HttpServer> clearServers() {
started.set(false);
Collection<io.vertx.core.http.HttpServer> servers = httpServers;
httpServers = null;
return servers == null ? Collections.emptyList() : servers;
}

private SeparatedCharSequence parsePath(String url) {
while (url.charAt(url.length() - 1) == '/') {
url = url.substring(0, url.length() - 1);
Expand Down Expand Up @@ -302,8 +331,9 @@ public NetworkType getType() {

@Override
public void shutdown() {
if (httpServers != null) {
for (io.vertx.core.http.HttpServer httpServer : httpServers) {
Collection<io.vertx.core.http.HttpServer> servers = clearServers();
if (!servers.isEmpty()) {
for (io.vertx.core.http.HttpServer httpServer : servers) {
httpServer.close(res -> {
if (res.failed()) {
log.error(res.cause().getMessage(), res.cause());
Expand All @@ -312,14 +342,16 @@ public void shutdown() {
}
});
}
httpServers.clear();
httpServers = null;
}
}

@Override
public boolean isAlive() {
return httpServers != null && !httpServers.isEmpty();
Collection<io.vertx.core.http.HttpServer> servers = httpServers;
return started.get() &&
servers != null &&
!servers.isEmpty() &&
servers.stream().allMatch(server -> server.actualPort() > 0);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
import org.jetlinks.community.network.security.VertxKeyCertTrustOptions;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

import javax.annotation.Nonnull;
Expand Down Expand Up @@ -75,30 +76,36 @@ public Mono<Network> createNetwork(@Nonnull VertxMqttServerProperties properties
private Mono<Network> initServer(VertxMqttServer server, VertxMqttServerProperties properties) {
int numberOfInstance = Math.max(1, properties.getInstance());
return convert(properties)
.map(options -> {
.flatMap(options -> {
List<MqttServer> instances = new ArrayList<>(numberOfInstance);
for (int i = 0; i < numberOfInstance; i++) {
MqttServer mqttServer = MqttServer.create(vertx, options);
instances.add(mqttServer);
}
server.setBind(new InetSocketAddress(options.getHost(), options.getPort()));
server.setLastError(null);
server.setMqttServer(instances);
for (MqttServer instance : instances) {
vertx.nettyEventLoopGroup()
.execute(()->{
instance.listen(result -> {
if (result.succeeded()) {
log.debug("startup mqtt server [{}] on port :{} ", properties.getId(), result
.result()
.actualPort());
} else {
server.setLastError(result.cause().getMessage());
log.warn("startup mqtt server [{}] error ", properties.getId(), result.cause());
}
});
});
}
return server;
return Flux
.fromIterable(instances)
.flatMap(instance -> Mono
.fromCompletionStage(instance.listen().toCompletionStage())
.doOnNext(result -> log.debug(
"startup mqtt server [{}] on port :{} ",
properties.getId(),
result.actualPort()
)))
.doOnCancel(server::shutdown)
.then(Mono.fromSupplier(() -> {
server.startupComplete();
return server;
}))
.onErrorResume(error -> {
server.setLastError(error.getMessage());
log.warn("startup mqtt server [{}] error ", properties.getId(), error);
return server
.shutdownAsync()
.then(Mono.error(error));
});
});

}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,15 +27,18 @@
import org.jetlinks.community.network.DefaultNetworkType;
import org.jetlinks.community.network.NetworkType;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.core.publisher.Sinks;
import reactor.util.concurrent.Queues;

import java.net.InetSocketAddress;
import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.atomic.AtomicBoolean;

@Slf4j
public class VertxMqttServer implements MqttServer {
Expand All @@ -45,7 +48,9 @@ public class VertxMqttServer implements MqttServer {
private final Map<String, List<Sinks.Many<MqttConnection>>> sinks =
new NonBlockingHashMap<>();

private Collection<io.vertx.mqtt.MqttServer> mqttServer;
private volatile Collection<io.vertx.mqtt.MqttServer> mqttServer;

private final AtomicBoolean started = new AtomicBoolean();

private final String id;

Expand All @@ -61,6 +66,7 @@ public VertxMqttServer(String id) {
}

public void setMqttServer(Collection<io.vertx.mqtt.MqttServer> mqttServer) {
started.set(false);
if (this.mqttServer != null && !this.mqttServer.isEmpty()) {
shutdown();
}
Expand All @@ -76,6 +82,29 @@ public void setMqttServer(Collection<io.vertx.mqtt.MqttServer> mqttServer) {
}
}

void startupComplete() {
started.set(true);
}

Mono<Void> shutdownAsync() {
return Flux
.fromIterable(clearServers())
.flatMap(server -> Mono
.fromCompletionStage(server.close().toCompletionStage())
.onErrorResume(error -> {
log.warn("close mqtt server error", error);
return Mono.empty();
}))
.then();
}

private synchronized Collection<io.vertx.mqtt.MqttServer> clearServers() {
started.set(false);
Collection<io.vertx.mqtt.MqttServer> servers = mqttServer;
mqttServer = null;
return servers == null ? Collections.emptyList() : servers;
}

private boolean emitNext(Sinks.Many<MqttConnection> sink, VertxMqttConnection connection){
if (sink.currentSubscriberCount() <= 0) {
return false;
Expand Down Expand Up @@ -130,7 +159,11 @@ public Flux<MqttConnection> handleConnection(String holder) {

@Override
public boolean isAlive() {
return mqttServer != null && !mqttServer.isEmpty();
Collection<io.vertx.mqtt.MqttServer> servers = mqttServer;
return started.get() &&
servers != null &&
!servers.isEmpty() &&
servers.stream().allMatch(server -> server.actualPort() > 0);
}

@Override
Expand All @@ -150,8 +183,9 @@ public NetworkType getType() {

@Override
public void shutdown() {
if (mqttServer != null) {
for (io.vertx.mqtt.MqttServer server : mqttServer) {
Collection<io.vertx.mqtt.MqttServer> servers = clearServers();
if (!servers.isEmpty()) {
for (io.vertx.mqtt.MqttServer server : servers) {
server.close(res -> {
if (res.failed()) {
log.error(res.cause().getMessage(), res.cause());
Expand All @@ -160,7 +194,6 @@ public void shutdown() {
}
});
}
mqttServer.clear();
}

}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import org.jetlinks.community.network.tcp.parser.PayloadParserBuilder;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

import javax.annotation.Nonnull;
Expand Down Expand Up @@ -88,7 +89,7 @@ public Mono<Network> createNetwork(@Nonnull TcpServerProperties properties) {

private Mono<Network> initTcpServer(VertxTcpServer tcpServer, TcpServerProperties properties) {
return convert(properties)
.map(options -> {
.flatMap(options -> {
int instance = Math.max(2, properties.getInstance());
List<NetServer> instances = new ArrayList<>(instance);
for (int i = 0; i < instance; i++) {
Expand All @@ -97,26 +98,32 @@ private Mono<Network> initTcpServer(VertxTcpServer tcpServer, TcpServerPropertie
Supplier<PayloadParser> parserSupplier= payloadParserBuilder.build(properties.getParserType(), properties);
parserSupplier.get();

tcpServer.setLastError(null);
tcpServer.setParserSupplier(parserSupplier);
tcpServer.setServer(instances);
tcpServer.setKeepAliveTimeout(properties.getLong("keepAliveTimeout", Duration
.ofMinutes(10)
.toMillis()));
tcpServer.setBind(new InetSocketAddress(properties.getHost(), properties.getPort()));
for (NetServer netServer : instances) {
vertx.nettyEventLoopGroup()
.execute(()->{
netServer.listen(properties.createSocketAddress(), result -> {
if (result.succeeded()) {
log.info("tcp server startup on {}", result.result().actualPort());
} else {
tcpServer.setLastError(result.cause().getMessage());
log.error("startup tcp server error", result.cause());
}
});
});
}
return tcpServer;
return Flux
.fromIterable(instances)
.flatMap(netServer -> Mono
.fromCompletionStage(netServer
.listen(properties.createSocketAddress())
.toCompletionStage())
.doOnNext(server -> log.info("tcp server startup on {}", server.actualPort())))
.doOnCancel(tcpServer::shutdown)
.then(Mono.fromSupplier(() -> {
tcpServer.startupComplete();
return tcpServer;
}))
.onErrorResume(error -> {
tcpServer.setLastError(error.getMessage());
log.error("startup tcp server error", error);
return tcpServer
.shutdownAsync()
.then(Mono.error(error));
});
});
}

Expand Down
Loading