fix storage corruption because of thread's interrupted state

This commit is contained in:
Eugene Zhuravlev
2012-02-29 17:25:04 +04:00
parent 5a9e5f5acc
commit 2ad4e229d1
6 changed files with 86 additions and 58 deletions
@@ -323,7 +323,7 @@ public class CompileServerManager implements ApplicationComponent{
synchronized (myAutomakeFutures) {
for (Map.Entry<RequestFuture, Project> entry : myAutomakeFutures.entrySet()) {
if (entry.getValue().equals(project)) {
entry.getKey().cancel(true);
entry.getKey().cancel(false);
}
}
}
@@ -609,7 +609,7 @@ public class CompileDriver {
if (future != null) {
while (!future.waitFor(200L , TimeUnit.MILLISECONDS)) {
if (indicator.isCanceled()) {
future.cancel(true);
future.cancel(false);
}
}
}
@@ -108,7 +108,7 @@ public class CompileServerClient extends SimpleProtobufClient<JpsServerResponseH
protected void beforeDisconnect() {
final ScheduledFuture<?> future = myPingFuture;
if (future != null) {
future.cancel(true);
future.cancel(false);
myPingFuture = null;
}
}
@@ -372,7 +372,7 @@ public class JavaBuilder extends ModuleLevelBuilder {
);
while (!future.waitFor(100L, TimeUnit.MILLISECONDS)) {
if (context.isCanceled()) {
future.cancel(true);
future.cancel(false);
}
}
rc = future.getResponseHandler().isTerminatedSuccessfully();
@@ -17,6 +17,7 @@ import org.jboss.netty.handler.codec.protobuf.ProtobufVarint32FrameDecoder;
import org.jboss.netty.handler.codec.protobuf.ProtobufVarint32LengthFieldPrepender;
import org.jetbrains.annotations.NonNls;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.jps.api.AsyncTaskExecutor;
import org.jetbrains.jps.api.GlobalOptions;
import org.jetbrains.jps.api.JpsRemoteProto;
import org.jetbrains.jps.incremental.Paths;
@@ -55,7 +56,22 @@ public class Server {
myBuildsExecutor = Executors.newFixedThreadPool(MAX_SIMULTANEOUS_BUILD_SESSIONS);
myChannelFactory = new NioServerSocketChannelFactory(threadPool, threadPool, 1);
final ChannelRegistrar channelRegistrar = new ChannelRegistrar();
myMessageHandler = new ServerMessageHandler(myBuildsExecutor, this);
myMessageHandler = new ServerMessageHandler(this, new AsyncTaskExecutor() {
@Override
public void submit(final Runnable runnable) {
myBuildsExecutor.submit(new Runnable() {
@Override
public void run() {
try {
runnable.run();
}
finally {
Thread.interrupted(); // clear interrupted status before returning to pull
}
}
});
}
});
myPipelineFactory = new ChannelPipelineFactory() {
public ChannelPipeline getPipeline() throws Exception {
return Channels.pipeline(
@@ -114,14 +130,12 @@ public class Server {
}
private void doStop(long elapsedTime) {
if (!myMessageHandler.hasRunningBuilds()) {
try {
System.out.println("Stopping compile server; reason: no pings from client received in " + elapsedTime + " ms");
stop();
}
finally {
System.exit(0);
}
try {
System.out.println("Stopping compile server; reason: no pings from client received in " + elapsedTime + " ms");
myMessageHandler.cancelAllBuildsAndClearState();
}
finally {
stop();
}
}
}, allowedIdlePeriod, allowedIdlePeriod, TimeUnit.MILLISECONDS);
@@ -129,8 +143,8 @@ public class Server {
public void stop() {
try {
myScheduler.shutdownNow();
myBuildsExecutor.shutdownNow();
myScheduler.shutdown();
myBuildsExecutor.shutdown();
final ChannelGroupFuture closeFuture = myAllOpenChannels.close();
closeFuture.awaitUninterruptibly();
}
@@ -162,7 +176,12 @@ public class Server {
final Server server = new Server(systemDir);
Runtime.getRuntime().addShutdownHook(new Thread("Shutdown hook thread") {
public void run() {
server.stop();
try {
server.myMessageHandler.cancelAllBuildsAndClearState();
}
finally {
server.stop();
}
}
});
@@ -15,7 +15,6 @@ import java.io.File;
import java.io.PrintStream;
import java.util.*;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.RunnableFuture;
/**
@@ -27,12 +26,12 @@ class ServerMessageHandler extends SimpleChannelHandler {
private final Map<String, SequentialTaskExecutor> myTaskExecutors = new HashMap<String, SequentialTaskExecutor>();
private final List<Pair<RunnableFuture, CompilationTask>> myBuildsInProgress = Collections.synchronizedList(new LinkedList<Pair<RunnableFuture, CompilationTask>>());
private final ExecutorService myBuildsExecutor;
private final Server myServer;
private final AsyncTaskExecutor myAsyncExecutor;
public ServerMessageHandler(ExecutorService buildsExecutor, Server server) {
myBuildsExecutor = buildsExecutor;
public ServerMessageHandler(Server server, final AsyncTaskExecutor asyncExecutor) {
myServer = server;
myAsyncExecutor = asyncExecutor;
}
public void messageReceived(final ChannelHandlerContext ctx, MessageEvent e) throws Exception {
@@ -87,35 +86,14 @@ class ServerMessageHandler extends SimpleChannelHandler {
break;
case SHUTDOWN_COMMAND :
myBuildsExecutor.submit(new Runnable() {
myAsyncExecutor.submit(new Runnable() {
public void run() {
final List<RunnableFuture> futures = new ArrayList<RunnableFuture>();
synchronized (myBuildsInProgress) {
for (Iterator<Pair<RunnableFuture, CompilationTask>> it = myBuildsInProgress.iterator(); it.hasNext(); ) {
final Pair<RunnableFuture, CompilationTask> pair = it.next();
it.remove();
pair.second.cancel();
final RunnableFuture future = pair.first;
futures.add(future);
future.cancel(true);
}
try {
cancelAllBuildsAndClearState();
}
facade.clearCahedState();
// wait until really stopped
for (RunnableFuture future : futures) {
try {
future.get();
}
catch (InterruptedException ignored) {
}
catch (ExecutionException ignored) {
}
finally {
myServer.stop();
}
myServer.stop();
}
});
break;
@@ -124,16 +102,24 @@ class ServerMessageHandler extends SimpleChannelHandler {
final String projectId = fsEvent.getProjectId();
final ProjectDescriptor pd = facade.getProjectDescriptor(projectId);
if (pd != null) {
final boolean wasInterrupted = Thread.interrupted();
try {
for (String path : fsEvent.getChangedPathsList()) {
facade.notifyFileChanged(pd, new File(path));
try {
for (String path : fsEvent.getChangedPathsList()) {
facade.notifyFileChanged(pd, new File(path));
}
for (String path : fsEvent.getDeletedPathsList()) {
facade.notifyFileDeleted(pd, new File(path));
}
}
for (String path : fsEvent.getDeletedPathsList()) {
facade.notifyFileDeleted(pd, new File(path));
finally {
pd.release();
}
}
finally {
pd.release();
if (wasInterrupted) {
Thread.currentThread().interrupt();
}
}
}
reply = ProtoUtil.toMessage(sessionId, ProtoUtil.createCommandCompletedEvent(null));
@@ -149,6 +135,34 @@ class ServerMessageHandler extends SimpleChannelHandler {
}
}
public void cancelAllBuildsAndClearState() {
final List<RunnableFuture> futures = new ArrayList<RunnableFuture>();
synchronized (myBuildsInProgress) {
for (Iterator<Pair<RunnableFuture, CompilationTask>> it = myBuildsInProgress.iterator(); it.hasNext(); ) {
final Pair<RunnableFuture, CompilationTask> pair = it.next();
it.remove();
pair.second.cancel();
final RunnableFuture future = pair.first;
futures.add(future);
future.cancel(false);
}
}
ServerState.getInstance().clearCahedState();
// wait until really stopped
for (RunnableFuture future : futures) {
try {
future.get();
}
catch (InterruptedException ignored) {
}
catch (ExecutionException ignored) {
}
}
}
private void cancelSession(UUID targetSessionId) {
synchronized (myBuildsInProgress) {
for (Iterator<Pair<RunnableFuture, CompilationTask>> it = myBuildsInProgress.iterator(); it.hasNext(); ) {
@@ -157,7 +171,7 @@ class ServerMessageHandler extends SimpleChannelHandler {
if (task.getSessionId().equals(targetSessionId)) {
it.remove();
task.cancel();
pair.first.cancel(true);
pair.first.cancel(false);
break;
}
}
@@ -212,12 +226,7 @@ class ServerMessageHandler extends SimpleChannelHandler {
synchronized (myTaskExecutors) {
SequentialTaskExecutor executor = myTaskExecutors.get(projectId);
if (executor == null) {
executor = new SequentialTaskExecutor(new AsyncTaskExecutor() {
@Override
public void submit(Runnable runnable) {
myBuildsExecutor.submit(runnable);
}
});
executor = new SequentialTaskExecutor(myAsyncExecutor);
myTaskExecutors.put(projectId, executor);
}
return executor;