Promise.blockingGet, PasswordSafe.getAsync

This commit is contained in:
Vladimir Krivosheev
2016-09-16 10:24:59 +02:00
parent d6d4db8e34
commit 95284119fd
14 changed files with 144 additions and 118 deletions
@@ -15,10 +15,7 @@ import com.intellij.openapi.util.Getter
import com.intellij.openapi.util.Key
import com.intellij.util.Consumer
import com.intellij.util.net.NetUtils
import org.jetbrains.concurrency.AsyncPromise
import org.jetbrains.concurrency.AsyncValueLoader
import org.jetbrains.concurrency.Promise
import org.jetbrains.concurrency.isRejected
import org.jetbrains.concurrency.*
import javax.swing.Icon
abstract class NetService @JvmOverloads protected constructor(protected val project: Project, private val consoleManager: ConsoleManager = ConsoleManager()) : Disposable {
@@ -45,7 +42,7 @@ abstract class NetService @JvmOverloads protected constructor(protected val proj
promise.rejected {
processHandler.destroyProcess()
Promise.logError(LOG, it)
LOG.errorIfNotMessage(it)
}
val processListener = MyProcessAdapter(processHandler)
@@ -28,6 +28,7 @@ import io.netty.handler.codec.http.*
import org.jetbrains.builtInWebServer.SingleConnectionNetService
import org.jetbrains.concurrency.Promise
import org.jetbrains.concurrency.doneRun
import org.jetbrains.concurrency.errorIfNotMessage
import org.jetbrains.io.*
import java.util.concurrent.atomic.AtomicInteger
@@ -87,7 +88,7 @@ abstract class FastCgiService(project: Project) : SingleConnectionNetService(pro
promise
.doneRun { fastCgiRequest.writeToServerChannel(notEmptyContent, processChannel.get()!!) }
.rejected {
Promise.logError(LOG, it)
LOG.errorIfNotMessage(it)
handleError(fastCgiRequest, notEmptyContent)
}
}
@@ -24,6 +24,7 @@ import com.intellij.ide.passwordSafe.PasswordStorage
import com.intellij.openapi.application.ApplicationManager
import com.intellij.openapi.components.SettingsSavingComponent
import com.intellij.openapi.diagnostic.catchAndLog
import org.jetbrains.concurrency.runAsync
class PasswordSafeImpl(/* public - backward compatibility */val settings: PasswordSafeSettings) : PasswordSafe(), SettingsSavingComponent {
private @Volatile var currentProvider: PasswordStorage
@@ -95,6 +96,9 @@ class PasswordSafeImpl(/* public - backward compatibility */val settings: Passwo
}
}
// maybe in the future we will use native async, so, this method added here instead "if need, just use runAsync in your code"
override fun getAsync(attributes: CredentialAttributes) = runAsync { get(attributes) }
override fun save() {
(currentProvider as? KeePassCredentialStore)?.let { it.save() }
}
@@ -20,6 +20,7 @@ import com.intellij.credentialStore.Credentials;
import com.intellij.openapi.components.ServiceManager;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.concurrency.Promise;
public abstract class PasswordSafe implements PasswordStorage {
@NotNull
@@ -30,4 +31,7 @@ public abstract class PasswordSafe implements PasswordStorage {
public abstract void set(@NotNull CredentialAttributes attributes, @Nullable Credentials credentials, boolean memoryOnly);
public abstract boolean isMemoryOnly();
@NotNull
public abstract Promise<Credentials> getAsync(@NotNull CredentialAttributes attributes);
}
@@ -30,6 +30,7 @@ import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.concurrency.AsyncPromise;
import org.jetbrains.concurrency.Promise;
import org.jetbrains.concurrency.PromiseKt;
import java.io.File;
import java.io.IOException;
@@ -305,7 +306,7 @@ public class RemoteFileInfoImpl implements RemoteContentProvider.DownloadingCall
case ERROR_OCCURRED:
default:
return Promise.reject("errorOccured");
return PromiseKt.rejectedPromise("errorOccurred");
}
}
}
@@ -19,7 +19,11 @@ import com.intellij.openapi.diagnostic.Logger
import com.intellij.openapi.util.Getter
import com.intellij.util.Consumer
import com.intellij.util.Function
import org.jetbrains.concurrency.Promise.State
import java.util.*
import java.util.concurrent.CountDownLatch
import java.util.concurrent.TimeUnit
import java.util.concurrent.TimeoutException
import java.util.concurrent.atomic.AtomicReference
private val LOG = Logger.getInstance(AsyncPromise::class.java)
@@ -27,7 +31,7 @@ private val LOG = Logger.getInstance(AsyncPromise::class.java)
@SuppressWarnings("ThrowableResultOfMethodCallIgnored")
private val OBSOLETE_ERROR = Promise.createError("Obsolete")
open class AsyncPromise<T> : Promise<T>(), Getter<T> {
open class AsyncPromise<T> : Promise<T>, Getter<T> {
private val doneRef = AtomicReference<Consumer<in T>?>()
private val rejectedRef = AtomicReference<Consumer<in Throwable>?>()
@@ -195,7 +199,7 @@ open class AsyncPromise<T> : Promise<T>(), Getter<T> {
doneRef.set(null)
if (rejected == null) {
Promise.logError(LOG, error)
LOG.errorIfNotMessage(error)
}
else if (!isObsolete(rejected)) {
rejected.consume(error)
@@ -218,6 +222,22 @@ open class AsyncPromise<T> : Promise<T>(), Getter<T> {
return this
}
override fun blockingGet(timeout: Int, timeUnit: TimeUnit): T? {
val latch = CountDownLatch(1)
processed { latch.countDown() }
if (!latch.await(timeout.toLong(), timeUnit)) {
throw TimeoutException()
}
@Suppress("UNCHECKED_CAST")
if (isRejected) {
throw (result as Throwable)
}
else {
return result as T?
}
}
private fun <T> setHandler(ref: AtomicReference<Consumer<in T>?>, newConsumer: Consumer<in T>, targetState: State) {
while (true) {
val oldConsumer = ref.get()
@@ -313,6 +333,4 @@ inline fun <T> AsyncPromise<*>.catchError(runnable: () -> T): T? {
private val cancelledPromise = RejectedPromise<Any?>(OBSOLETE_ERROR)
@Suppress("UNCHECKED_CAST")
fun <T> cancelledPromise(): Promise<T> = cancelledPromise as Promise<T>
fun <T> rejectedPromise(error: Throwable): Promise<T> = Promise.reject(error)
fun <T> cancelledPromise(): Promise<T> = cancelledPromise as Promise<T>
@@ -19,8 +19,11 @@ import com.intellij.openapi.util.Getter;
import com.intellij.util.Consumer;
import com.intellij.util.Function;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
class DonePromise<T> extends Promise<T> implements Getter<T> {
import java.util.concurrent.TimeUnit;
class DonePromise<T> implements Getter<T>, Promise<T> {
private final T result;
public DonePromise(T result) {
@@ -59,7 +62,7 @@ class DonePromise<T> extends Promise<T> implements Getter<T> {
@Override
public <SUB_RESULT> Promise<SUB_RESULT> then(@NotNull Function<? super T, ? extends SUB_RESULT> done) {
if (done instanceof Obsolescent && ((Obsolescent)done).isObsolete()) {
return Promise.reject("obsolete");
return PromiseKt.rejectedPromise("obsolete");
}
else {
return Promise.resolve(done.fun(result));
@@ -78,6 +81,12 @@ class DonePromise<T> extends Promise<T> implements Getter<T> {
return State.FULFILLED;
}
@Nullable
@Override
public T blockingGet(int timeout, @NotNull TimeUnit timeUnit) {
return result;
}
@Override
public T get() {
return result;
@@ -15,32 +15,30 @@
*/
package org.jetbrains.concurrency;
import com.intellij.openapi.application.ApplicationManager;
import com.intellij.openapi.diagnostic.Logger;
import com.intellij.openapi.progress.ProcessCanceledException;
import com.intellij.openapi.util.ActionCallback;
import com.intellij.openapi.util.AsyncResult;
import com.intellij.util.Consumer;
import com.intellij.util.Function;
import com.intellij.util.ThreeState;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
public abstract class Promise<T> {
public static final Promise<Void> DONE = new DonePromise<>(null);
public static final Promise<Void> REJECTED = new RejectedPromise<>(createError("rejected"));
import java.util.concurrent.TimeUnit;
public interface Promise<T> {
Promise<Void> DONE = new DonePromise<>(null);
Promise<Void> REJECTED = PromiseKt.getREJECTED();
@NotNull
public static RuntimeException createError(@NotNull String error) {
static RuntimeException createError(@NotNull String error) {
return new MessageError(error);
}
public enum State {
enum State {
PENDING, FULFILLED, REJECTED
}
@NotNull
public static <T> Promise<T> resolve(T result) {
static <T> Promise<T> resolve(T result) {
if (result == null) {
//noinspection unchecked
return (Promise<T>)DONE;
@@ -51,23 +49,7 @@ public abstract class Promise<T> {
}
@NotNull
public static <T> Promise<T> reject(@NotNull String error) {
return reject(createError(error));
}
@NotNull
public static <T> Promise<T> reject(@Nullable Throwable error) {
if (error == null) {
//noinspection unchecked
return (Promise<T>)REJECTED;
}
else {
return new RejectedPromise<>(error);
}
}
@NotNull
public static Promise<Void> wrapAsVoid(@NotNull ActionCallback asyncResult) {
static Promise<Void> wrapAsVoid(@NotNull ActionCallback asyncResult) {
final AsyncPromise<Void> promise = new AsyncPromise<>();
asyncResult.doWhenDone(() -> promise.setResult(null)).doWhenRejected(
error -> promise.setError(createError(error == null ? "Internal error" : error)));
@@ -75,78 +57,34 @@ public abstract class Promise<T> {
}
@NotNull
public static <T> Promise<T> wrap(@NotNull AsyncResult<T> asyncResult) {
static <T> Promise<T> wrap(@NotNull AsyncResult<T> asyncResult) {
final AsyncPromise<T> promise = new AsyncPromise<>();
asyncResult.doWhenDone(new Consumer<T>() {
@Override
public void consume(T result) {
promise.setResult(result);
}
}).doWhenRejected(promise::setError);
asyncResult.doWhenDone((Consumer<T>)result -> promise.setResult(result)).doWhenRejected(promise::setError);
return promise;
}
@NotNull
public abstract Promise<T> done(@NotNull Consumer<? super T> done);
Promise<T> done(@NotNull Consumer<? super T> done);
@NotNull
public abstract Promise<T> processed(@NotNull AsyncPromise<? super T> fulfilled);
Promise<T> processed(@NotNull AsyncPromise<? super T> fulfilled);
@NotNull
public abstract Promise<T> rejected(@NotNull Consumer<Throwable> rejected);
Promise<T> rejected(@NotNull Consumer<Throwable> rejected);
public abstract Promise<T> processed(@NotNull Consumer<? super T> processed);
Promise<T> processed(@NotNull Consumer<? super T> processed);
@NotNull
public abstract <SUB_RESULT> Promise<SUB_RESULT> then(@NotNull Function<? super T, ? extends SUB_RESULT> done);
<SUB_RESULT> Promise<SUB_RESULT> then(@NotNull Function<? super T, ? extends SUB_RESULT> done);
@NotNull
public abstract <SUB_RESULT> Promise<SUB_RESULT> thenAsync(@NotNull Function<? super T, Promise<SUB_RESULT>> done);
<SUB_RESULT> Promise<SUB_RESULT> thenAsync(@NotNull Function<? super T, Promise<SUB_RESULT>> done);
@NotNull
public abstract State getState();
State getState();
@SuppressWarnings("ExceptionClassNameDoesntEndWithException")
public static class MessageError extends RuntimeException {
private final ThreeState log;
@Nullable
T blockingGet(int timeout, @NotNull TimeUnit timeUnit);
public MessageError(@NotNull String error) {
super(error);
log = ThreeState.UNSURE;
}
public MessageError(@NotNull String error, boolean log) {
super(error);
this.log = ThreeState.fromBoolean(log);
}
@NotNull
@Override
public final synchronized Throwable fillInStackTrace() {
return this;
}
}
/**
* Log error if not message error
*/
public static boolean logError(@NotNull Logger logger, @NotNull Throwable e) {
if (e instanceof MessageError) {
ThreeState log = ((MessageError)e).log;
if (log == ThreeState.YES || (log == ThreeState.UNSURE && ApplicationManager.getApplication().isUnitTestMode())) {
logger.error(e);
return true;
}
}
else if (!(e instanceof ProcessCanceledException)) {
logger.error(e);
return true;
}
return false;
}
public abstract void notify(@NotNull AsyncPromise<? super T> child);
void notify(@NotNull AsyncPromise<? super T> child);
}
@@ -20,7 +20,7 @@ import com.intellij.util.Function
import java.util.concurrent.TimeUnit
internal class RejectedPromise<T>(private val error: Throwable) : Promise<T>() {
internal class RejectedPromise<T>(private val error: Throwable) : Promise<T> {
override fun getState() = Promise.State.REJECTED
override fun done(done: Consumer<in T>) = this
@@ -15,14 +15,16 @@
*/
package org.jetbrains.concurrency
import com.intellij.openapi.application.ApplicationManager
import com.intellij.openapi.diagnostic.Logger
import com.intellij.openapi.progress.ProcessCanceledException
import com.intellij.util.Consumer
import com.intellij.util.Function
import com.intellij.util.SmartList
import com.intellij.util.ThreeState
import com.intellij.util.concurrency.AppExecutorUtil
import java.util.*
private val rejectedPromise = Promise.reject<Any?>("rejected")
// only internal usage
interface ObsolescentFunction<Param, Result> : Function<Param, Result>, Obsolescent
@@ -80,11 +82,10 @@ inline fun Promise<*>.rejected(node: Obsolescent, crossinline handler: (Throwabl
override fun consume(param: Throwable) = handler(param)
})
fun <T> rejectedPromise(error: String): Promise<T> = Promise.reject(error)
val REJECTED: Promise<Void> = RejectedPromise(createError("rejected", false))
@Suppress("UNCHECKED_CAST")
fun <T> rejectedPromise(): Promise<T> = rejectedPromise as Promise<T>
fun <T> rejectedPromise(): Promise<T> = REJECTED as Promise<T>
val Promise<*>.isRejected: Boolean
get() = state == Promise.State.REJECTED
@@ -107,7 +108,7 @@ fun <T> collectResults(promises: List<Promise<T>>): Promise<List<T>> {
return all(promises, results)
}
fun createError(error: String, log: Boolean): RuntimeException = Promise.MessageError(error, log)
fun createError(error: String, log: Boolean): RuntimeException = MessageError(error, log)
inline fun <T> AsyncPromise<T>.compute(runnable: () -> T) {
val result = catchError(runnable)
@@ -116,11 +117,66 @@ inline fun <T> AsyncPromise<T>.compute(runnable: () -> T) {
}
}
inline fun runAsync(crossinline runnable: () -> Unit): Promise<*> {
val promise = AsyncPromise<Any?>()
inline fun <T> runAsync(crossinline runnable: () -> T): Promise<T> {
val promise = AsyncPromise<T>()
AppExecutorUtil.getAppExecutorService().execute {
promise.catchError { runnable() }
promise.setResult(null)
val result = try {
runnable()
}
catch (e: Throwable) {
promise.setError(e)
return@execute
}
promise.setResult(result)
}
return promise
}
fun <T> rejectedPromise(error: String): Promise<T> = rejectedPromise(createError(error, true))
fun <T> rejectedPromise(error: Throwable?): Promise<T> {
if (error == null) {
@Suppress("UNCHECKED_CAST")
return REJECTED as Promise<T>
}
else {
return RejectedPromise(error)
}
}
@SuppressWarnings("ExceptionClassNameDoesntEndWithException")
internal class MessageError : RuntimeException {
internal val log: ThreeState
constructor(error: String) : super(error) {
log = ThreeState.UNSURE
}
constructor(error: String, log: Boolean) : super(error) {
this.log = ThreeState.fromBoolean(log)
}
fun fillInStackTrace(): Throwable {
return this
}
}
/**
* Log error if not a message error
*/
fun Logger.errorIfNotMessage(e: Throwable): Boolean {
if (e is MessageError) {
val log = e.log
if (log == ThreeState.YES || (log == ThreeState.UNSURE && (ApplicationManager.getApplication()?.isUnitTestMode ?: false))) {
error(e)
return true
}
}
else if (e !is ProcessCanceledException) {
error(e)
return true
}
return false
}
@@ -20,6 +20,7 @@ import com.intellij.util.io.shutdownIfOio
import io.netty.channel.Channel
import org.jetbrains.concurrency.AsyncPromise
import org.jetbrains.concurrency.Promise
import org.jetbrains.concurrency.errorIfNotMessage
import org.jetbrains.concurrency.resolvedPromise
import org.jetbrains.jsonProtocol.Request
import org.jetbrains.rpc.CONNECTION_CLOSED_MESSAGE
@@ -61,7 +62,7 @@ open class StandaloneVmHelper(private val vm: Vm, private val messageProcessor:
messageProcessor.send(disconnectRequest)
.rejected {
if (it.message != CONNECTION_CLOSED_MESSAGE) {
Promise.logError(LOG, it)
LOG.errorIfNotMessage(it)
}
}
// we don't wait response because 1) no response to "disconnect" message (V8 for example) 2) closed message manager just ignore any incoming messages
@@ -18,8 +18,8 @@ package org.jetbrains.debugger
import com.intellij.xdebugger.frame.XCompositeNode
import com.intellij.xdebugger.frame.XValueChildrenList
import com.intellij.xdebugger.frame.XValueGroup
import org.jetbrains.concurrency.Promise
import org.jetbrains.concurrency.done
import org.jetbrains.concurrency.errorIfNotMessage
import org.jetbrains.debugger.values.FunctionValue
import org.jetbrains.rpc.LOG
import java.util.*
@@ -39,7 +39,7 @@ internal class FunctionScopesValueGroup(private val functionValue: FunctionValue
}
}
.rejected {
Promise.logError(LOG, it)
LOG.errorIfNotMessage(it)
node.setErrorMessage(it.message!!)
}
}
@@ -26,10 +26,7 @@ import com.intellij.util.io.socketConnection.ConnectionStatus
import io.netty.bootstrap.Bootstrap
import io.netty.channel.ChannelFuture
import io.netty.util.concurrent.GenericFutureListener
import org.jetbrains.concurrency.AsyncPromise
import org.jetbrains.concurrency.Promise
import org.jetbrains.concurrency.rejectedPromise
import org.jetbrains.concurrency.resolvedPromise
import org.jetbrains.concurrency.*
import org.jetbrains.debugger.Vm
import org.jetbrains.io.NettyUtil
import org.jetbrains.rpc.LOG
@@ -76,7 +73,7 @@ abstract class RemoteVmConnection : VmConnection<Vm>() {
}
.rejected {
if (it !is ConnectException) {
Promise.logError(LOG, it)
LOG.errorIfNotMessage(it)
}
setState(ConnectionStatus.CONNECTION_FAILED, it.message)
}
@@ -17,12 +17,12 @@ package org.jetbrains.debugger
import com.intellij.util.Consumer
import com.intellij.xdebugger.XDebugSession
import org.jetbrains.concurrency.Promise
import org.jetbrains.concurrency.errorIfNotMessage
import org.jetbrains.rpc.LOG
class RejectErrorReporter @JvmOverloads constructor(private val session: XDebugSession, private val description: String? = null) : Consumer<Throwable> {
override fun consume(error: Throwable) {
if (Promise.logError(LOG, error)) {
if (LOG.errorIfNotMessage(error)) {
session.reportError("${if (description == null) "" else "$description: "}${error.message}")
}
}