initial pub serve proxy (based on Alexander Doroshko patch)

This commit is contained in:
Vladimir Krivosheev
2014-09-12 20:28:11 +02:00
parent 25198531c7
commit c5bb65a7c7
10 changed files with 286 additions and 195 deletions
@@ -19,6 +19,7 @@ import com.intellij.openapi.diagnostic.Logger;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelHandlerAdapter;
import io.netty.channel.ChannelHandlerContext;
import org.jetbrains.annotations.NotNull;
import java.net.ConnectException;
@@ -31,6 +32,7 @@ public final class ChannelExceptionHandler extends ChannelHandlerAdapter {
private ChannelExceptionHandler() {
}
@NotNull
public static ChannelHandler getInstance() {
return INSTANCE;
}
@@ -63,12 +63,13 @@ public final class NettyUtil {
}
}
public static Channel connectClient(Bootstrap bootstrap, InetSocketAddress remoteAddress, ActionCallback asyncResult) {
@Nullable
public static Channel connectClient(@NotNull Bootstrap bootstrap, @NotNull InetSocketAddress remoteAddress, @Nullable ActionCallback asyncResult) {
return connect(bootstrap, remoteAddress, asyncResult, DEFAULT_CONNECT_ATTEMPT_COUNT);
}
@Nullable
public static Channel connect(@NotNull Bootstrap bootstrap, @NotNull InetSocketAddress remoteAddress, @NotNull ActionCallback asyncResult, int maxAttemptCount) {
public static Channel connect(@NotNull Bootstrap bootstrap, @NotNull InetSocketAddress remoteAddress, @Nullable ActionCallback asyncResult, int maxAttemptCount) {
try {
int attemptCount = 0;
@@ -85,7 +86,9 @@ public final class NettyUtil {
else {
@SuppressWarnings("ThrowableResultOfMethodCallIgnored")
Throwable cause = future.cause();
asyncResult.reject("Cannot connect: " + (cause == null ? "unknown error" : cause.getMessage()));
if (asyncResult != null) {
asyncResult.reject("Cannot connect: " + (cause == null ? "unknown error" : cause.getMessage()));
}
return null;
}
}
@@ -104,7 +107,9 @@ public final class NettyUtil {
Thread.sleep(attemptCount * MIN_START_TIME);
}
else {
asyncResult.reject("Cannot connect: " + e.getMessage());
if (asyncResult != null) {
asyncResult.reject("Cannot connect: " + e.getMessage());
}
return null;
}
}
@@ -115,7 +120,9 @@ public final class NettyUtil {
return channel;
}
catch (Throwable e) {
asyncResult.reject("Cannot connect: " + e.getMessage());
if (asyncResult != null) {
asyncResult.reject("Cannot connect: " + e.getMessage());
}
return null;
}
}
@@ -110,7 +110,7 @@ public final class BuiltInWebServer extends HttpRequestHandler {
else {
projectName = host;
}
return doProcess(request, context.channel(), projectName);
return doProcess(request, context, projectName);
}
public static boolean isOwnHostName(@NotNull String host) {
@@ -135,7 +135,7 @@ public final class BuiltInWebServer extends HttpRequestHandler {
}
}
private static boolean doProcess(@NotNull FullHttpRequest request, @NotNull Channel channel, @Nullable String projectName) {
private static boolean doProcess(@NotNull FullHttpRequest request, @NotNull ChannelHandlerContext context, @Nullable String projectName) {
final String decodedPath = URLUtil.unescapePercentSequences(UriUtil.trimParameters(request.uri()));
int offset;
boolean emptyPath;
@@ -163,7 +163,7 @@ public final class BuiltInWebServer extends HttpRequestHandler {
}
// we must redirect "jsdebug" to "jsdebug/" as nginx does, otherwise browser will treat it as file instead of directory, so, relative path will not work
WebServerPathHandler.redirectToDirectory(request, channel, projectName);
WebServerPathHandler.redirectToDirectory(request, context.channel(), projectName);
return true;
}
@@ -172,7 +172,7 @@ public final class BuiltInWebServer extends HttpRequestHandler {
for (WebServerPathHandler pathHandler : WebServerPathHandler.EP_NAME.getExtensions()) {
try {
if (pathHandler.process(path, project, request, channel, projectName, decodedPath, isCustomHost)) {
if (pathHandler.process(path, project, request, context, projectName, decodedPath, isCustomHost)) {
return true;
}
}
@@ -18,6 +18,7 @@ package org.jetbrains.builtInWebServer;
import com.intellij.openapi.project.Project;
import com.intellij.openapi.vfs.VirtualFile;
import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.http.FullHttpRequest;
import io.netty.handler.codec.http.HttpResponseStatus;
import org.jetbrains.annotations.NotNull;
@@ -29,10 +30,11 @@ final class DefaultWebServerPathHandler extends WebServerPathHandler {
public boolean process(@NotNull String path,
@NotNull Project project,
@NotNull FullHttpRequest request,
@NotNull Channel channel,
@NotNull ChannelHandlerContext context,
@Nullable String projectName,
@NotNull String decodedRawPath,
boolean isCustomHost) {
Channel channel = context.channel();
WebServerPathToFileManager pathToFileManager = WebServerPathToFileManager.getInstance(project);
VirtualFile result = pathToFileManager.pathToFileCache.getIfPresent(path);
boolean indexUsed = false;
@@ -0,0 +1,183 @@
package org.jetbrains.builtInWebServer;
import com.intellij.concurrency.JobScheduler;
import com.intellij.execution.ExecutionException;
import com.intellij.execution.filters.TextConsoleBuilder;
import com.intellij.execution.filters.TextConsoleBuilderFactory;
import com.intellij.execution.process.OSProcessHandler;
import com.intellij.execution.process.ProcessAdapter;
import com.intellij.execution.process.ProcessEvent;
import com.intellij.execution.ui.ConsoleView;
import com.intellij.execution.ui.ConsoleViewContentType;
import com.intellij.openapi.Disposable;
import com.intellij.openapi.application.ApplicationManager;
import com.intellij.openapi.diagnostic.Logger;
import com.intellij.openapi.project.Project;
import com.intellij.openapi.util.AsyncResult;
import com.intellij.openapi.util.AsyncValueLoader;
import com.intellij.openapi.util.Key;
import com.intellij.openapi.wm.ToolWindow;
import com.intellij.openapi.wm.ToolWindowAnchor;
import com.intellij.openapi.wm.ToolWindowManager;
import com.intellij.ui.content.ContentFactory;
import com.intellij.util.Consumer;
import com.intellij.util.net.NetUtils;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.io.NettyUtil;
import javax.swing.*;
import java.io.IOException;
import java.util.concurrent.TimeUnit;
public abstract class NetService implements Disposable {
protected static final Logger LOG = Logger.getInstance(NetService.class);
protected final Project project;
protected final AsyncValueLoader<OSProcessHandler> processHandler = new AsyncValueLoader<OSProcessHandler>() {
@Override
protected boolean isCancelOnReject() {
return true;
}
@Nullable
private OSProcessHandler doGetProcessHandler(int port) {
try {
return createProcessHandler(project, port);
}
catch (ExecutionException e) {
LOG.error(e);
return null;
}
}
@Override
protected void load(@NotNull final AsyncResult<OSProcessHandler> result) throws IOException {
final int port = NetUtils.findAvailableSocketPort();
final OSProcessHandler processHandler = doGetProcessHandler(port);
if (processHandler == null) {
result.setRejected();
return;
}
result.doWhenRejected(new Runnable() {
@Override
public void run() {
processHandler.destroyProcess();
}
});
final MyProcessAdapter processListener = new MyProcessAdapter();
processHandler.addProcessListener(processListener);
processHandler.startNotify();
if (result.isRejected()) {
return;
}
JobScheduler.getScheduler().schedule(new Runnable() {
@Override
public void run() {
if (result.isRejected()) {
return;
}
ApplicationManager.getApplication().executeOnPooledThread(new Runnable() {
@Override
public void run() {
if (!result.isRejected()) {
try {
connectToProcess(result, port, processHandler, processListener);
}
catch (Throwable e) {
result.setRejected();
LOG.error(e);
}
}
}
});
}
}, NettyUtil.MIN_START_TIME, TimeUnit.MILLISECONDS);
}
@Override
protected void disposeResult(@NotNull OSProcessHandler processHandler) {
try {
closeProcessConnections();
}
finally {
processHandler.destroyProcess();
}
}
};
private ConsoleView console;
protected NetService(@NotNull Project project) {
this.project = project;
}
@Nullable
protected abstract OSProcessHandler createProcessHandler(Project project, int port) throws ExecutionException;
protected void connectToProcess(@NotNull AsyncResult<OSProcessHandler> asyncResult, int port, @NotNull OSProcessHandler processHandler, @NotNull Consumer<String> errorOutputConsumer) {
asyncResult.setDone(processHandler);
}
protected abstract void closeProcessConnections();
@Override
public void dispose() {
processHandler.reset();
}
protected void configureConsole(@NotNull TextConsoleBuilder consoleBuilder) {
}
@NotNull
protected abstract String getConsoleToolWindowId();
@NotNull
protected abstract Icon getConsoleToolWindowIcon();
private final class MyProcessAdapter extends ProcessAdapter implements Consumer<String> {
private void createConsole() {
TextConsoleBuilder consoleBuilder = TextConsoleBuilderFactory.getInstance().createBuilder(project);
configureConsole(consoleBuilder);
console = consoleBuilder.getConsole();
ApplicationManager.getApplication().invokeLater(new Runnable() {
@Override
public void run() {
ToolWindow toolWindow = ToolWindowManager.getInstance(project).registerToolWindow(getConsoleToolWindowId(), false, ToolWindowAnchor.BOTTOM, project, true);
toolWindow.setIcon(getConsoleToolWindowIcon());
toolWindow.getContentManager().addContent(ContentFactory.SERVICE.getInstance().createContent(console.getComponent(), "", false));
}
}, project.getDisposed());
}
@Override
public void onTextAvailable(ProcessEvent event, Key outputType) {
print(event.getText(), ConsoleViewContentType.getConsoleViewType(outputType));
}
private void print(String text, ConsoleViewContentType contentType) {
if (console == null) {
createConsole();
}
console.print(text, contentType);
}
@Override
public void processTerminated(ProcessEvent event) {
processHandler.reset();
print(getConsoleToolWindowId() + " terminated\n", ConsoleViewContentType.SYSTEM_OUTPUT);
}
@Override
public void consume(String message) {
print(message, ConsoleViewContentType.ERROR_OUTPUT);
}
}
}
@@ -0,0 +1,42 @@
package org.jetbrains.builtInWebServer;
import com.intellij.execution.process.OSProcessHandler;
import com.intellij.openapi.project.Project;
import com.intellij.openapi.util.AsyncResult;
import com.intellij.util.Consumer;
import com.intellij.util.net.NetUtils;
import io.netty.bootstrap.Bootstrap;
import io.netty.channel.Channel;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.io.NettyUtil;
import java.net.InetSocketAddress;
public abstract class SingleConnectionNetService extends NetService {
protected volatile Channel processChannel;
protected SingleConnectionNetService(@NotNull Project project) {
super(project);
}
protected abstract void configureBootstrap(@NotNull Bootstrap bootstrap, @NotNull Consumer<String> errorOutputConsumer);
@Override
protected void connectToProcess(@NotNull AsyncResult<OSProcessHandler> asyncResult, int port, @NotNull OSProcessHandler processHandler, @NotNull Consumer<String> errorOutputConsumer) {
Bootstrap bootstrap = NettyUtil.oioClientBootstrap();
configureBootstrap(bootstrap, errorOutputConsumer);
processChannel = NettyUtil.connectClient(bootstrap, new InetSocketAddress(NetUtils.getLoopbackAddress(), port), asyncResult);
if (processChannel != null) {
asyncResult.setDone(processHandler);
}
}
@Override
protected void closeProcessConnections() {
Channel currentProcessChannel = processChannel;
if (currentProcessChannel != null) {
processChannel = null;
NettyUtil.closeAndReleaseFactory(currentProcessChannel);
}
}
}
@@ -19,6 +19,7 @@ import com.intellij.openapi.extensions.ExtensionPointName;
import com.intellij.openapi.project.Project;
import com.intellij.openapi.vfs.VfsUtil;
import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.http.*;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
@@ -38,7 +39,7 @@ public abstract class WebServerPathHandler {
public abstract boolean process(@NotNull String path,
@NotNull Project project,
@NotNull FullHttpRequest request,
@NotNull Channel channel,
@NotNull ChannelHandlerContext context,
@Nullable String projectName,
@NotNull String decodedRawPath,
boolean isCustomHost);
@@ -16,22 +16,22 @@
package org.jetbrains.builtInWebServer;
import com.intellij.openapi.project.Project;
import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.http.FullHttpRequest;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
public abstract class WebServerPathHandlerAdapter extends WebServerPathHandler {
protected abstract boolean process(@NotNull String path, @NotNull Project project, @NotNull FullHttpRequest request, @NotNull Channel channel);
protected abstract boolean process(@NotNull String path, @NotNull Project project, @NotNull FullHttpRequest request, @NotNull ChannelHandlerContext context);
@Override
public final boolean process(@NotNull String path,
@NotNull Project project,
@NotNull FullHttpRequest request,
@NotNull Channel channel,
@NotNull ChannelHandlerContext context,
@Nullable String projectName,
@NotNull String decodedRawPath,
boolean isCustomHost) {
return process(path, project, request, channel);
return process(path, project, request, context);
}
}
@@ -8,6 +8,7 @@ import io.netty.channel.Channel;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.http.*;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.io.Responses;
import org.jetbrains.io.SimpleChannelInboundHandlerAdapter;
@@ -17,7 +18,7 @@ import static org.jetbrains.io.fastCgi.FastCgiService.LOG;
public class FastCgiChannelHandler extends SimpleChannelInboundHandlerAdapter<FastCgiResponse> {
private final ConcurrentIntObjectMap<Channel> requestToChannel;
public FastCgiChannelHandler(ConcurrentIntObjectMap<Channel> channel) {
public FastCgiChannelHandler(@NotNull ConcurrentIntObjectMap<Channel> channel) {
requestToChannel = channel;
}
@@ -1,151 +1,61 @@
package org.jetbrains.io.fastCgi;
import com.intellij.concurrency.JobScheduler;
import com.intellij.execution.filters.TextConsoleBuilder;
import com.intellij.execution.filters.TextConsoleBuilderFactory;
import com.intellij.execution.process.OSProcessHandler;
import com.intellij.execution.process.ProcessAdapter;
import com.intellij.execution.process.ProcessEvent;
import com.intellij.execution.ui.ConsoleView;
import com.intellij.execution.ui.ConsoleViewContentType;
import com.intellij.openapi.Disposable;
import com.intellij.openapi.application.ApplicationManager;
import com.intellij.openapi.diagnostic.Logger;
import com.intellij.openapi.project.Project;
import com.intellij.openapi.util.AsyncResult;
import com.intellij.openapi.util.AsyncValueLoader;
import com.intellij.openapi.util.Key;
import com.intellij.openapi.wm.ToolWindow;
import com.intellij.openapi.wm.ToolWindowAnchor;
import com.intellij.openapi.wm.ToolWindowManager;
import com.intellij.ui.content.ContentFactory;
import com.intellij.util.Consumer;
import com.intellij.util.containers.ContainerUtil;
import com.intellij.util.containers.StripedLockIntObjectConcurrentHashMap;
import com.intellij.util.net.NetUtils;
import io.netty.bootstrap.Bootstrap;
import io.netty.buffer.ByteBuf;
import io.netty.channel.Channel;
import io.netty.channel.ChannelInitializer;
import io.netty.handler.codec.http.HttpResponseStatus;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.builtInWebServer.SingleConnectionNetService;
import org.jetbrains.io.ChannelExceptionHandler;
import org.jetbrains.io.NettyUtil;
import org.jetbrains.io.Responses;
import javax.swing.*;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
// todo send FCGI_ABORT_REQUEST if client channel disconnected
public abstract class FastCgiService implements Disposable {
public abstract class FastCgiService extends SingleConnectionNetService {
static final Logger LOG = Logger.getInstance(FastCgiService.class);
protected final Project project;
private final AtomicInteger requestIdCounter = new AtomicInteger();
private final StripedLockIntObjectConcurrentHashMap<Channel> requests = new StripedLockIntObjectConcurrentHashMap<Channel>();
protected final StripedLockIntObjectConcurrentHashMap<Channel> requests = new StripedLockIntObjectConcurrentHashMap<Channel>();
private volatile Channel fastCgiChannel;
protected final AsyncValueLoader<OSProcessHandler> processHandler = new AsyncValueLoader<OSProcessHandler>() {
@Override
protected boolean isCancelOnReject() {
return true;
}
@Override
protected void load(@NotNull final AsyncResult<OSProcessHandler> result) throws IOException {
final int port = NetUtils.findAvailableSocketPort();
final OSProcessHandler processHandler = createProcessHandler(project, port);
if (processHandler == null) {
result.setRejected();
return;
}
result.doWhenRejected(new Runnable() {
@Override
public void run() {
processHandler.destroyProcess();
}
});
final MyProcessAdapter processListener = new MyProcessAdapter();
processHandler.addProcessListener(processListener);
processHandler.startNotify();
if (result.isRejected()) {
return;
}
JobScheduler.getScheduler().schedule(new Runnable() {
@Override
public void run() {
if (result.isRejected()) {
return;
}
ApplicationManager.getApplication().executeOnPooledThread(new Runnable() {
@Override
public void run() {
if (!result.isRejected()) {
try {
connectToProcess(result, port, processHandler, processListener);
}
catch (Throwable e) {
result.setRejected();
LOG.error(e);
}
}
}
});
}
}, NettyUtil.MIN_START_TIME, TimeUnit.MILLISECONDS);
}
@Override
protected void disposeResult(@NotNull OSProcessHandler processHandler) {
try {
Channel currentFastCgiChannel = fastCgiChannel;
if (currentFastCgiChannel != null) {
fastCgiChannel = null;
NettyUtil.closeAndReleaseFactory(currentFastCgiChannel);
}
processHandler.destroyProcess();
}
finally {
requestIdCounter.set(0);
if (!requests.isEmpty()) {
List<Channel> waitingClients = ContainerUtil.toList(requests.elements());
requests.clear();
for (Channel channel : waitingClients) {
try {
if (channel.isActive()) {
Responses.sendStatus(HttpResponseStatus.BAD_GATEWAY, channel);
}
}
catch (Throwable e) {
NettyUtil.log(e, LOG);
}
}
}
}
}
};
private ConsoleView console;
protected FastCgiService(Project project) {
this.project = project;
public FastCgiService(@NotNull Project project) {
super(project);
}
protected abstract OSProcessHandler createProcessHandler(Project project, int port);
@Override
protected void closeProcessConnections() {
try {
super.closeProcessConnections();
}
finally {
requestIdCounter.set(0);
if (!requests.isEmpty()) {
List<Channel> waitingClients = ContainerUtil.toList(requests.elements());
requests.clear();
for (Channel channel : waitingClients) {
try {
if (channel.isActive()) {
Responses.sendStatus(HttpResponseStatus.BAD_GATEWAY, channel);
}
}
catch (Throwable e) {
NettyUtil.log(e, LOG);
}
}
}
}
}
private void connectToProcess(final AsyncResult<OSProcessHandler> asyncResult, final int port, final OSProcessHandler processHandler, final Consumer<String> errorOutputConsumer) {
Bootstrap bootstrap = NettyUtil.oioClientBootstrap();
@Override
protected void configureBootstrap(@NotNull Bootstrap bootstrap, @NotNull final Consumer<String> errorOutputConsumer) {
final FastCgiChannelHandler fastCgiChannelHandler = new FastCgiChannelHandler(requests);
bootstrap.handler(new ChannelInitializer() {
@Override
@@ -153,23 +63,19 @@ public abstract class FastCgiService implements Disposable {
channel.pipeline().addLast(new FastCgiDecoder(errorOutputConsumer), fastCgiChannelHandler, ChannelExceptionHandler.getInstance());
}
});
fastCgiChannel = NettyUtil.connectClient(bootstrap, new InetSocketAddress(NetUtils.getLoopbackAddress(), port), asyncResult);
if (fastCgiChannel != null) {
asyncResult.setDone(processHandler);
}
}
public void send(final FastCgiRequest fastCgiRequest, final ByteBuf content) {
content.retain();
if (processHandler.has()) {
fastCgiRequest.writeToServerChannel(content, fastCgiChannel);
fastCgiRequest.writeToServerChannel(content, processChannel);
}
else {
processHandler.get().doWhenDone(new Runnable() {
@Override
public void run() {
fastCgiRequest.writeToServerChannel(content, fastCgiChannel);
fastCgiRequest.writeToServerChannel(content, processChannel);
}
}).doWhenRejected(new Runnable() {
@Override
@@ -193,57 +99,4 @@ public abstract class FastCgiService implements Disposable {
requests.put(requestId, channel);
return requestId;
}
@Override
public void dispose() {
processHandler.reset();
}
protected abstract void buildConsole(@NotNull TextConsoleBuilder consoleBuilder);
@NotNull
protected abstract String getConsoleToolWindowId();
@NotNull
protected abstract Icon getConsoleToolWindowIcon();
private final class MyProcessAdapter extends ProcessAdapter implements Consumer<String> {
private void createConsole() {
TextConsoleBuilder consoleBuilder = TextConsoleBuilderFactory.getInstance().createBuilder(project);
buildConsole(consoleBuilder);
console = consoleBuilder.getConsole();
ApplicationManager.getApplication().invokeLater(new Runnable() {
@Override
public void run() {
ToolWindow toolWindow = ToolWindowManager.getInstance(project).registerToolWindow(getConsoleToolWindowId(), false, ToolWindowAnchor.BOTTOM, project, true);
toolWindow.setIcon(getConsoleToolWindowIcon());
toolWindow.getContentManager().addContent(ContentFactory.SERVICE.getInstance().createContent(console.getComponent(), "", false));
}
}, project.getDisposed());
}
@Override
public void onTextAvailable(ProcessEvent event, Key outputType) {
print(event.getText(), ConsoleViewContentType.getConsoleViewType(outputType));
}
private void print(String text, ConsoleViewContentType contentType) {
if (console == null) {
createConsole();
}
console.print(text, contentType);
}
@Override
public void processTerminated(ProcessEvent event) {
processHandler.reset();
print(getConsoleToolWindowId() + " terminated\n", ConsoleViewContentType.SYSTEM_OUTPUT);
}
@Override
public void consume(String message) {
print(message, ConsoleViewContentType.ERROR_OUTPUT);
}
}
}