GH-750 - introduce dedicated thread pool executors to avoid handing threads to block everything else

This commit is contained in:
Martin Lippert
2022-04-13 16:42:07 +02:00
parent cda0b75048
commit 5d0bed47e4
3 changed files with 29 additions and 17 deletions

View File

@@ -17,6 +17,8 @@ import java.util.concurrent.CancellationException;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.Executor;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
import java.util.stream.Collectors;
@@ -90,6 +92,8 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
private final ListenerList<TextDocument> documentCloseListeners = new ListenerList<>();
private final ListenerList<TextDocument> documentOpenListeners = new ListenerList<>();
private List<Consumer<TextDocumentSaveChange>> documentSaveListeners = ImmutableList.of();
private final Executor messageWorkerThreadPool;
private CompletionHandler completionHandler;
private CompletionResolveHandler completionResolveHandler;
@@ -104,6 +108,8 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
public SimpleTextDocumentService(SimpleLanguageServer server, LanguageServerProperties props) {
this.server = server;
this.props = props;
this.messageWorkerThreadPool = Executors.newCachedThreadPool();
}
/**
@@ -253,7 +259,7 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
public CompletableFuture<Either<List<CompletionItem>, CompletionList>> completion(CompletionParams position) {
log.info("completion request arrived: " + position.getTextDocument().getUri());
return CompletableFutures.computeAsync(cancelToken -> {
return CompletableFutures.computeAsync(messageWorkerThreadPool, cancelToken -> {
CompletionHandler h = completionHandler;
if (h != null) {
@@ -269,7 +275,7 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
public CompletableFuture<CompletionItem> resolveCompletionItem(CompletionItem unresolved) {
log.info("Completion item resolve request received: {}", unresolved.getLabel());
return CompletableFutures.computeAsync(cancelToken -> {
return CompletableFutures.computeAsync(messageWorkerThreadPool, cancelToken -> {
try {
CompletionResolveHandler h = completionResolveHandler;
if (h != null) {
@@ -291,7 +297,7 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
public CompletableFuture<Hover> hover(HoverParams hoverParams) {
log.debug("hover requested for {}", hoverParams.getPosition());
CompletableFuture<Hover> result = CompletableFutures.computeAsync(cancelToken -> {
CompletableFuture<Hover> result = CompletableFutures.computeAsync(messageWorkerThreadPool, cancelToken -> {
return computeHover(cancelToken, hoverParams);
});
@@ -330,7 +336,7 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
DefinitionHandler h = this.definitionHandler;
if (h != null) {
return CompletableFutures.computeAsync(cancelToken -> {
return CompletableFutures.computeAsync(messageWorkerThreadPool, cancelToken -> {
cancelToken.checkCanceled();
@@ -361,7 +367,7 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
ReferencesHandler h = this.referencesHandler;
if (h != null) {
return CompletableFutures.computeAsync(cancelToken -> {
return CompletableFutures.computeAsync(messageWorkerThreadPool, cancelToken -> {
List<? extends Location> list = h.handle(cancelToken, params);
return list != null && list.isEmpty() ? null : list;
});
@@ -376,7 +382,7 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
DocumentSymbolHandler h = this.documentSymbolHandler;
if (h != null) {
return CompletableFutures.computeAsync(cancelToken -> {
return CompletableFutures.computeAsync(messageWorkerThreadPool, cancelToken -> {
cancelToken.checkCanceled();
try {
@@ -433,7 +439,7 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
CodeLensHandler handler = this.codeLensHandler;
if (handler != null) {
return CompletableFutures.computeAsync(cancelToken -> {
return CompletableFutures.computeAsync(messageWorkerThreadPool, cancelToken -> {
return handler.handle(cancelToken, params);
});
}
@@ -445,7 +451,7 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
CodeLensResolveHandler handler = this.codeLensResolveHandler;
if (handler != null) {
return CompletableFutures.computeAsync(cancelToken -> {
return CompletableFutures.computeAsync(messageWorkerThreadPool, cancelToken -> {
return handler.handle(unresolved);
});
@@ -476,7 +482,7 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
}
}
}
});
}, messageWorkerThreadPool);
}
}
@@ -484,7 +490,7 @@ public class SimpleTextDocumentService implements TextDocumentService, DocumentE
public CompletableFuture<List<? extends DocumentHighlight>> documentHighlight(DocumentHighlightParams highlightParams) {
DocumentHighlightHandler handler = this.documentHighlightHandler;
if (handler != null) {
return CompletableFutures.computeAsync(cancelToken -> {
return CompletableFutures.computeAsync(messageWorkerThreadPool, cancelToken -> {
return handler.handle(cancelToken, highlightParams);
});

View File

@@ -1,5 +1,5 @@
/*******************************************************************************
* Copyright (c) 2019, 2020 Pivotal, Inc.
* Copyright (c) 2019, 2022 Pivotal, Inc.
* All rights reserved. This program and the accompanying materials
* are made available under the terms of the Eclipse Public License v1.0
* which accompanies this distribution, and is available at
@@ -20,6 +20,8 @@ import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executor;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import org.slf4j.Logger;
@@ -36,7 +38,6 @@ import com.sun.tools.attach.VirtualMachineDescriptor;
*/
@SuppressWarnings("restriction")
public class SpringProcessConnectorLocal {
private static final Logger log = LoggerFactory.getLogger(SpringProcessConnectorLocal.class);
@@ -47,12 +48,16 @@ public class SpringProcessConnectorLocal {
private final Set<SpringProcessDescriptor> processes;
private final SpringProcessConnectorService processConnectorService;
private final Executor statusUpdateThreadPool;
private boolean projectsChanged;
public SpringProcessConnectorLocal(SpringProcessConnectorService processConnector, ProjectObserver projectObserver) {
this.projects = new ConcurrentHashMap<>();
this.processes = Collections.synchronizedSet(new HashSet<>());
this.statusUpdateThreadPool = Executors.newFixedThreadPool(10);
this.projectsChanged = false;
this.processConnectorService = processConnector;
@@ -168,7 +173,7 @@ public class SpringProcessConnectorLocal {
List<CompletableFuture<Void>> futures = new ArrayList<>();
for (SpringProcessDescriptor process : processes) {
futures.add(process.updateStatus(projects::containsKey, projects::get));
futures.add(process.updateStatus(projects::containsKey, projects::get, statusUpdateThreadPool));
}
CompletableFuture<Void> allStatusUpdates = CompletableFuture.allOf((CompletableFuture[]) futures.toArray(new CompletableFuture[futures.size()]));

View File

@@ -1,5 +1,5 @@
/*******************************************************************************
* Copyright (c) 2019, 2020 Pivotal, Inc.
* Copyright (c) 2019, 2022 Pivotal, Inc.
* All rights reserved. This program and the accompanying materials
* are made available under the terms of the Eclipse Public License v1.0
* which accompanies this distribution, and is available at
@@ -12,6 +12,7 @@ package org.springframework.ide.vscode.boot.java.livehover.v2;
import java.util.Properties;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;
import java.util.function.Predicate;
import org.slf4j.Logger;
@@ -57,7 +58,7 @@ public class SpringProcessDescriptor {
private SpringProcessStatus status;
private String projectName;
public SpringProcessDescriptor(String processKey, String processID, String processName, VirtualMachineDescriptor vm) {
this.processKey = processKey;
this.processID = processID;
@@ -109,11 +110,11 @@ public class SpringProcessDescriptor {
return true;
}
public CompletableFuture<Void> updateStatus(Predicate<String> projectIsKnown, Predicate<String> projectHasActuators) {
public CompletableFuture<Void> updateStatus(Predicate<String> projectIsKnown, Predicate<String> projectHasActuators, Executor statusUpdateThreadPool) {
return CompletableFuture.supplyAsync(() -> {
this.status = checkStatus(projectIsKnown, projectHasActuators);
return null;
});
}, statusUpdateThreadPool);
}
private SpringProcessStatus checkStatus(Predicate<String> projectIsKnown, Predicate<String> projectHasActuators) {