Fix TODO in ReactorHttpServer
This commit is contained in:
@@ -16,6 +16,8 @@
|
|||||||
|
|
||||||
package org.springframework.http.server.reactive.bootstrap;
|
package org.springframework.http.server.reactive.bootstrap;
|
||||||
|
|
||||||
|
import java.util.concurrent.atomic.AtomicReference;
|
||||||
|
|
||||||
import reactor.core.Loopback;
|
import reactor.core.Loopback;
|
||||||
import reactor.ipc.netty.NettyContext;
|
import reactor.ipc.netty.NettyContext;
|
||||||
|
|
||||||
@@ -31,7 +33,7 @@ public class ReactorHttpServer extends HttpServerSupport implements HttpServer,
|
|||||||
|
|
||||||
private reactor.ipc.netty.http.server.HttpServer reactorServer;
|
private reactor.ipc.netty.http.server.HttpServer reactorServer;
|
||||||
|
|
||||||
private NettyContext running;
|
private AtomicReference<NettyContext> nettyContext = new AtomicReference<>();
|
||||||
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@@ -43,42 +45,39 @@ public class ReactorHttpServer extends HttpServerSupport implements HttpServer,
|
|||||||
Assert.notNull(getHttpHandler());
|
Assert.notNull(getHttpHandler());
|
||||||
this.reactorHandler = new ReactorHttpHandlerAdapter(getHttpHandler());
|
this.reactorHandler = new ReactorHttpHandlerAdapter(getHttpHandler());
|
||||||
}
|
}
|
||||||
this.reactorServer = reactor.ipc.netty.http.server.HttpServer.create(getHost(),
|
this.reactorServer = reactor.ipc.netty.http.server.HttpServer
|
||||||
getPort());
|
.create(getHost(), getPort());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public boolean isRunning() {
|
public boolean isRunning() {
|
||||||
NettyContext running = this.running;
|
NettyContext context = this.nettyContext.get();
|
||||||
return running != null && running.channel()
|
return (context != null && context.channel().isActive());
|
||||||
.isActive();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public Object connectedInput() {
|
public Object connectedInput() {
|
||||||
return reactorServer;
|
return this.reactorServer;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public Object connectedOutput() {
|
public Object connectedOutput() {
|
||||||
return reactorServer;
|
return this.reactorServer;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void start() {
|
public void start() {
|
||||||
// TODO: should be made thread-safe (compareAndSet..)
|
if (this.nettyContext.get() == null) {
|
||||||
if (this.running == null) {
|
this.nettyContext.set(this.reactorServer.newHandler(reactorHandler).block());
|
||||||
this.running = this.reactorServer.newHandler(reactorHandler).block();
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void stop() {
|
public void stop() {
|
||||||
NettyContext running = this.running;
|
NettyContext context = this.nettyContext.getAndSet(null);
|
||||||
if (running != null) {
|
if (context != null) {
|
||||||
this.running = null;
|
context.dispose();
|
||||||
running.dispose();
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user