mirror of
https://gitflic.ru/project/openide/openide.git
synced 2026-09-27 10:03:11 +07:00
convert ObjectValue, MessageManager, Value, CommandProcessor, ValueModifier, MessageWriter, ArrayValue, ValueBase, CommandSenderBase, ResultReader, MessageProcessor, ObjectValueBase, ValueNodeAsyncFunction, RequestCallback, MessageManagerBase to kotlin
This commit is contained in:
@@ -1,19 +1,17 @@
|
||||
package org.jetbrains.debugger;
|
||||
package org.jetbrains.debugger
|
||||
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.concurrency.Promise;
|
||||
import org.jetbrains.debugger.values.Value;
|
||||
import org.jetbrains.concurrency.Promise
|
||||
import org.jetbrains.debugger.values.Value
|
||||
|
||||
public interface ValueModifier {
|
||||
interface ValueModifier {
|
||||
// expression can contains reference to another variables in current scope, so, we should evaluate it before set
|
||||
// https://youtrack.jetbrains.com/issue/WEB-2342#comment=27-512122
|
||||
|
||||
// we don't worry about performance in case of simple primitive values - boolean/string/numbers,
|
||||
// it works quickly and we don't want to complicate our code and debugger SDK
|
||||
Promise<?> setValue(@NotNull Variable variable, String newValue, @NotNull EvaluateContext evaluateContext);
|
||||
fun setValue(variable: Variable, newValue: String, evaluateContext: EvaluateContext): Promise<*>
|
||||
|
||||
Promise<?> setValue(@NotNull Variable variable, @NotNull Value newValue, @NotNull EvaluateContext evaluateContext);
|
||||
fun setValue(variable: Variable, newValue: Value, evaluateContext: EvaluateContext): Promise<*>
|
||||
|
||||
@NotNull
|
||||
Promise<Value> evaluateGet(@NotNull Variable variable, @NotNull EvaluateContext evaluateContext);
|
||||
fun evaluateGet(variable: Variable, evaluateContext: EvaluateContext): Promise<Value>
|
||||
}
|
||||
@@ -1,14 +1,13 @@
|
||||
package org.jetbrains.debugger.values;
|
||||
package org.jetbrains.debugger.values
|
||||
|
||||
public interface ArrayValue extends ObjectValue {
|
||||
interface ArrayValue : ObjectValue {
|
||||
/**
|
||||
* Be aware - it is not equals to java array length.
|
||||
* In case of sparse array {@code
|
||||
* var sparseArray = [3, 4];
|
||||
* In case of sparse array `var sparseArray = [3, 4];
|
||||
* sparseArray[45] = 34;
|
||||
* sparseArray[40999995] = "foo";
|
||||
* }
|
||||
* sparseArray[40999995] = "foo";
|
||||
` *
|
||||
* length will be equal to 40999995.
|
||||
*/
|
||||
int getLength();
|
||||
val length: Int
|
||||
}
|
||||
+19
-29
@@ -1,49 +1,39 @@
|
||||
package org.jetbrains.debugger.values;
|
||||
package org.jetbrains.debugger.values
|
||||
|
||||
import com.intellij.util.ThreeState;
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.annotations.Nullable;
|
||||
import org.jetbrains.concurrency.Obsolescent;
|
||||
import org.jetbrains.concurrency.Promise;
|
||||
import org.jetbrains.debugger.EvaluateContext;
|
||||
import org.jetbrains.debugger.Variable;
|
||||
import org.jetbrains.debugger.VariablesHost;
|
||||
|
||||
import java.util.List;
|
||||
import com.intellij.util.ThreeState
|
||||
import org.jetbrains.concurrency.Obsolescent
|
||||
import org.jetbrains.concurrency.Promise
|
||||
import org.jetbrains.debugger.EvaluateContext
|
||||
import org.jetbrains.debugger.Variable
|
||||
import org.jetbrains.debugger.VariablesHost
|
||||
import org.jetbrains.debugger.Vm
|
||||
|
||||
/**
|
||||
* A compound value that has zero or more properties
|
||||
*/
|
||||
public interface ObjectValue extends Value {
|
||||
@Nullable
|
||||
String getClassName();
|
||||
interface ObjectValue : Value {
|
||||
val className: String?
|
||||
|
||||
@NotNull
|
||||
Promise<List<Variable>> getProperties();
|
||||
val properties: Promise<List<Variable>>
|
||||
|
||||
@NotNull
|
||||
Promise<List<Variable>> getProperties(@NotNull List<String> names, @NotNull EvaluateContext evaluateContext, @NotNull Obsolescent obsolescent);
|
||||
fun getProperties(names: List<String>, evaluateContext: EvaluateContext, obsolescent: Obsolescent): Promise<List<Variable>>
|
||||
|
||||
@NotNull
|
||||
VariablesHost getVariablesHost();
|
||||
val variablesHost: VariablesHost<ValueManager<Vm>>
|
||||
|
||||
/**
|
||||
* from (inclusive) to (exclusive) ranges of array elements or elements if less than bucketThreshold
|
||||
*
|
||||
|
||||
* "to" could be -1 (sometimes length is unknown, so, you can pass -1 instead of actual elements size)
|
||||
*/
|
||||
@NotNull
|
||||
Promise<Void> getIndexedProperties(int from, int to, int bucketThreshold, @NotNull IndexedVariablesConsumer consumer, @Nullable ValueType componentType);
|
||||
fun getIndexedProperties(from: Int, to: Int, bucketThreshold: Int, consumer: IndexedVariablesConsumer, componentType: ValueType?): Promise<*>
|
||||
|
||||
/**
|
||||
* It must return quickly. Return {@link com.intellij.util.ThreeState#UNSURE} otherwise.
|
||||
* It must return quickly. Return [com.intellij.util.ThreeState.UNSURE] otherwise.
|
||||
*/
|
||||
@NotNull
|
||||
ThreeState hasProperties();
|
||||
fun hasProperties() = ThreeState.UNSURE
|
||||
|
||||
/**
|
||||
* It must return quickly. Return {@link com.intellij.util.ThreeState#UNSURE} otherwise.
|
||||
* It must return quickly. Return [com.intellij.util.ThreeState.UNSURE] otherwise.
|
||||
*/
|
||||
@NotNull
|
||||
ThreeState hasIndexedProperties();
|
||||
fun hasIndexedProperties() = ThreeState.NO
|
||||
}
|
||||
+53
-107
@@ -1,127 +1,73 @@
|
||||
package org.jetbrains.debugger.values;
|
||||
package org.jetbrains.debugger.values
|
||||
|
||||
import com.intellij.util.SmartList;
|
||||
import com.intellij.util.ThreeState;
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.annotations.Nullable;
|
||||
import org.jetbrains.concurrency.Obsolescent;
|
||||
import org.jetbrains.concurrency.ObsolescentAsyncFunction;
|
||||
import org.jetbrains.concurrency.Promise;
|
||||
import org.jetbrains.debugger.EvaluateContext;
|
||||
import org.jetbrains.debugger.ValueModifier;
|
||||
import org.jetbrains.debugger.Variable;
|
||||
import org.jetbrains.debugger.VariablesHost;
|
||||
import com.intellij.util.SmartList
|
||||
import org.jetbrains.concurrency.Obsolescent
|
||||
import org.jetbrains.concurrency.ObsolescentAsyncFunction
|
||||
import org.jetbrains.concurrency.Promise
|
||||
import org.jetbrains.debugger.EvaluateContext
|
||||
import org.jetbrains.debugger.Variable
|
||||
import org.jetbrains.debugger.VariablesHost
|
||||
import org.jetbrains.debugger.Vm
|
||||
import java.util.*
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.Comparator;
|
||||
import java.util.List;
|
||||
abstract class ObjectValueBase<VALUE_LOADER : ValueManager<out Vm>>(type: ValueType) : ValueBase(type), ObjectValue {
|
||||
protected abstract val childrenManager: VariablesHost<VALUE_LOADER>
|
||||
|
||||
public abstract class ObjectValueBase<VALUE_LOADER extends ValueManager> extends ValueBase implements ObjectValue {
|
||||
protected VariablesHost<VALUE_LOADER> childrenManager;
|
||||
override val properties: Promise<List<Variable>>
|
||||
get() = childrenManager.get()
|
||||
|
||||
public ObjectValueBase(@NotNull ValueType type) {
|
||||
super(type);
|
||||
internal abstract inner class MyObsolescentAsyncFunction<PARAM, RESULT>(private val obsolescent: Obsolescent) : ObsolescentAsyncFunction<PARAM, RESULT> {
|
||||
override fun isObsolete() = obsolescent.isObsolete || childrenManager.valueManager.isObsolete
|
||||
}
|
||||
|
||||
@NotNull
|
||||
@Override
|
||||
public final Promise<List<Variable>> getProperties() {
|
||||
return childrenManager.get();
|
||||
}
|
||||
override fun getProperties(names: List<String>, evaluateContext: EvaluateContext, obsolescent: Obsolescent) = properties
|
||||
.then(object : MyObsolescentAsyncFunction<List<Variable>, List<Variable>>(obsolescent) {
|
||||
override fun `fun`(variables: List<Variable>) = getSpecifiedProperties(variables, names, evaluateContext)
|
||||
})
|
||||
|
||||
abstract class MyObsolescentAsyncFunction<PARAM, RESULT> implements ObsolescentAsyncFunction<PARAM, RESULT> {
|
||||
private final Obsolescent obsolescent;
|
||||
override val valueString: String? = null
|
||||
|
||||
MyObsolescentAsyncFunction(@NotNull Obsolescent obsolescent) {
|
||||
this.obsolescent = obsolescent;
|
||||
}
|
||||
override fun getIndexedProperties(from: Int, to: Int, bucketThreshold: Int, consumer: IndexedVariablesConsumer, componentType: ValueType?): Promise<*> = Promise.REJECTED
|
||||
|
||||
@Override
|
||||
public boolean isObsolete() {
|
||||
return obsolescent.isObsolete() || childrenManager.valueManager.isObsolete();
|
||||
}
|
||||
}
|
||||
@Suppress("CAST_NEVER_SUCCEEDS")
|
||||
override val variablesHost: VariablesHost<ValueManager<Vm>>
|
||||
get() = childrenManager as VariablesHost<ValueManager<Vm>>
|
||||
|
||||
@NotNull
|
||||
@Override
|
||||
public Promise<List<Variable>> getProperties(@NotNull final List<String> names, @NotNull final EvaluateContext evaluateContext, @NotNull final Obsolescent obsolescent) {
|
||||
return getProperties()
|
||||
.then(new MyObsolescentAsyncFunction<List<Variable>, List<Variable>>(obsolescent) {
|
||||
@NotNull
|
||||
@Override
|
||||
public Promise<List<Variable>> fun(List<Variable> variables) {
|
||||
return getSpecifiedProperties(variables, names, evaluateContext);
|
||||
companion object {
|
||||
protected fun getSpecifiedProperties(variables: List<Variable>, names: List<String>, evaluateContext: EvaluateContext): Promise<List<Variable>> {
|
||||
val properties = SmartList<Variable>()
|
||||
var getterCount = 0
|
||||
for (property in variables) {
|
||||
if (!property.isReadable || !names.contains(property.name)) {
|
||||
continue
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@NotNull
|
||||
protected static Promise<List<Variable>> getSpecifiedProperties(@NotNull List<Variable> variables, @NotNull final List<String> names, @NotNull EvaluateContext evaluateContext) {
|
||||
final List<Variable> properties = new SmartList<Variable>();
|
||||
int getterCount = 0;
|
||||
for (Variable property : variables) {
|
||||
if (!property.isReadable() || !names.contains(property.getName())) {
|
||||
continue;
|
||||
if (!properties.isEmpty()) {
|
||||
Collections.sort(properties, object : Comparator<Variable> {
|
||||
override fun compare(o1: Variable, o2: Variable) = names.indexOf(o1.name) - names.indexOf(o2.name)
|
||||
})
|
||||
}
|
||||
|
||||
properties.add(property)
|
||||
if (property.value == null) {
|
||||
getterCount++
|
||||
}
|
||||
}
|
||||
|
||||
if (!properties.isEmpty()) {
|
||||
Collections.sort(properties, new Comparator<Variable>() {
|
||||
@Override
|
||||
public int compare(@NotNull Variable o1, @NotNull Variable o2) {
|
||||
return names.indexOf(o1.getName()) - names.indexOf(o2.getName());
|
||||
if (getterCount == 0) {
|
||||
return Promise.resolve(properties)
|
||||
}
|
||||
else {
|
||||
val promises = SmartList<Promise<*>>()
|
||||
for (variable in properties) {
|
||||
if (variable.value == null) {
|
||||
val valueModifier = variable.valueModifier
|
||||
assert(valueModifier != null)
|
||||
promises.add(valueModifier!!.evaluateGet(variable, evaluateContext))
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
properties.add(property);
|
||||
if (property.getValue() == null) {
|
||||
getterCount++;
|
||||
}
|
||||
}
|
||||
|
||||
if (getterCount == 0) {
|
||||
return Promise.resolve(properties);
|
||||
}
|
||||
else {
|
||||
List<Promise<?>> promises = new SmartList<Promise<?>>();
|
||||
for (Variable variable : properties) {
|
||||
if (variable.getValue() == null) {
|
||||
ValueModifier valueModifier = variable.getValueModifier();
|
||||
assert valueModifier != null;
|
||||
promises.add(valueModifier.evaluateGet(variable, evaluateContext));
|
||||
}
|
||||
return Promise.all<List<Variable>>(promises, properties)
|
||||
}
|
||||
return Promise.all(promises, properties);
|
||||
}
|
||||
}
|
||||
|
||||
@Nullable
|
||||
@Override
|
||||
public String getValueString() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@NotNull
|
||||
@Override
|
||||
public ThreeState hasProperties() {
|
||||
return ThreeState.UNSURE;
|
||||
}
|
||||
|
||||
@NotNull
|
||||
@Override
|
||||
public ThreeState hasIndexedProperties() {
|
||||
return ThreeState.NO;
|
||||
}
|
||||
|
||||
@NotNull
|
||||
@Override
|
||||
public Promise<Void> getIndexedProperties(int from, int to, int bucketThreshold, @NotNull IndexedVariablesConsumer consumer, @Nullable ValueType componentType) {
|
||||
return Promise.REJECTED;
|
||||
}
|
||||
|
||||
@NotNull
|
||||
@Override
|
||||
public VariablesHost getVariablesHost() {
|
||||
return childrenManager;
|
||||
}
|
||||
}
|
||||
@@ -1,16 +1,13 @@
|
||||
package org.jetbrains.debugger.values;
|
||||
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
package org.jetbrains.debugger.values
|
||||
|
||||
/**
|
||||
* An object that represents a VM variable value (compound or atomic).
|
||||
*/
|
||||
public interface Value {
|
||||
@NotNull
|
||||
ValueType getType();
|
||||
interface Value {
|
||||
val type: ValueType
|
||||
|
||||
/**
|
||||
* @return a string representation of this value
|
||||
*/
|
||||
String getValueString();
|
||||
val valueString: String?
|
||||
}
|
||||
|
||||
@@ -1,17 +1,3 @@
|
||||
package org.jetbrains.debugger.values;
|
||||
package org.jetbrains.debugger.values
|
||||
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
|
||||
public abstract class ValueBase implements Value {
|
||||
protected final ValueType type;
|
||||
|
||||
public ValueBase(@NotNull ValueType type) {
|
||||
this.type = type;
|
||||
}
|
||||
|
||||
@NotNull
|
||||
@Override
|
||||
public final ValueType getType() {
|
||||
return type;
|
||||
}
|
||||
}
|
||||
abstract class ValueBase(override val type: ValueType) : Value
|
||||
@@ -1,44 +1,30 @@
|
||||
package org.jetbrains.rpc;
|
||||
package org.jetbrains.rpc
|
||||
|
||||
import com.intellij.openapi.diagnostic.Logger;
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.jsonProtocol.Request;
|
||||
import com.intellij.openapi.diagnostic.Logger
|
||||
import org.jetbrains.jsonProtocol.Request
|
||||
import java.util.concurrent.atomic.AtomicInteger
|
||||
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
val LOG = Logger.getInstance(CommandProcessor::class.java)
|
||||
|
||||
public abstract class CommandProcessor<INCOMING, INCOMING_WITH_SEQ, SUCCESS_RESPONSE>
|
||||
extends CommandSenderBase<SUCCESS_RESPONSE>
|
||||
implements MessageManager.Handler<Request, INCOMING, INCOMING_WITH_SEQ, SUCCESS_RESPONSE>,
|
||||
ResultReader<SUCCESS_RESPONSE>,
|
||||
MessageProcessor {
|
||||
public static final Logger LOG = Logger.getInstance(CommandProcessor.class);
|
||||
abstract class CommandProcessor<INCOMING, INCOMING_WITH_SEQ : Any, SUCCESS_RESPONSE>() : CommandSenderBase<SUCCESS_RESPONSE>(), MessageManager.Handler<Request<out Any>, INCOMING, INCOMING_WITH_SEQ, SUCCESS_RESPONSE>, ResultReader<SUCCESS_RESPONSE>, MessageProcessor {
|
||||
private val currentSequence = AtomicInteger()
|
||||
protected val messageManager = MessageManager(this)
|
||||
|
||||
private final AtomicInteger currentSequence = new AtomicInteger();
|
||||
protected final MessageManager<Request, INCOMING, INCOMING_WITH_SEQ, SUCCESS_RESPONSE> messageManager;
|
||||
|
||||
protected CommandProcessor() {
|
||||
messageManager = new MessageManager<Request, INCOMING, INCOMING_WITH_SEQ, SUCCESS_RESPONSE>(this);
|
||||
override fun cancelWaitingRequests() {
|
||||
messageManager.cancelWaitingRequests()
|
||||
}
|
||||
|
||||
@Override
|
||||
public final void cancelWaitingRequests() {
|
||||
messageManager.cancelWaitingRequests();
|
||||
override fun closed() {
|
||||
messageManager.closed()
|
||||
}
|
||||
|
||||
@Override
|
||||
public final void closed() {
|
||||
messageManager.closed();
|
||||
override fun getUpdatedSequence(message: Request<out Any>): Int {
|
||||
val id = currentSequence.incrementAndGet()
|
||||
message.finalize(id)
|
||||
return id
|
||||
}
|
||||
|
||||
@Override
|
||||
public final int getUpdatedSequence(@NotNull Request message) {
|
||||
int id = currentSequence.incrementAndGet();
|
||||
message.finalize(id);
|
||||
return id;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected <RESULT> void send(@NotNull Request message, @NotNull RequestPromise<SUCCESS_RESPONSE, RESULT> callback) {
|
||||
messageManager.send(message, callback);
|
||||
override final fun <RESULT : Any> doSend(message: Request<RESULT>, callback: CommandSenderBase.RequestPromise<SUCCESS_RESPONSE, RESULT>) {
|
||||
messageManager.send(message, callback)
|
||||
}
|
||||
}
|
||||
@@ -1,49 +1,43 @@
|
||||
package org.jetbrains.rpc;
|
||||
package org.jetbrains.rpc
|
||||
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.annotations.Nullable;
|
||||
import org.jetbrains.concurrency.AsyncPromise;
|
||||
import org.jetbrains.concurrency.Promise;
|
||||
import org.jetbrains.jsonProtocol.Request;
|
||||
import org.jetbrains.concurrency.AsyncPromise
|
||||
import org.jetbrains.concurrency.Promise
|
||||
import org.jetbrains.jsonProtocol.Request
|
||||
|
||||
public abstract class CommandSenderBase<SUCCESS_RESPONSE> implements CommandSender {
|
||||
protected abstract <RESULT> void send(@NotNull Request message, @NotNull RequestPromise<SUCCESS_RESPONSE, RESULT> callback);
|
||||
abstract class CommandSenderBase<SUCCESS_RESPONSE> {
|
||||
protected abstract fun <RESULT : Any> doSend(message: Request<RESULT>, callback: RequestPromise<SUCCESS_RESPONSE, RESULT>)
|
||||
|
||||
@Override
|
||||
@NotNull
|
||||
public final <RESULT> Promise<RESULT> send(@NotNull Request<RESULT> request) {
|
||||
RequestPromise<SUCCESS_RESPONSE, RESULT> callback = new RequestPromise<SUCCESS_RESPONSE, RESULT>(request.getMethodName());
|
||||
send(request, callback);
|
||||
return callback;
|
||||
fun <RESULT : Any> send(message: Request<RESULT>): Promise<RESULT> {
|
||||
val callback = RequestPromise<SUCCESS_RESPONSE, RESULT>(message.methodName)
|
||||
doSend(message, callback)
|
||||
return callback
|
||||
}
|
||||
|
||||
protected static final class RequestPromise<SUCCESS_RESPONSE, RESULT> extends AsyncPromise<RESULT> implements RequestCallback<SUCCESS_RESPONSE> {
|
||||
private final String methodName;
|
||||
|
||||
public RequestPromise(@Nullable String methodName) {
|
||||
this.methodName = methodName;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onSuccess(@Nullable SUCCESS_RESPONSE response, @Nullable ResultReader<SUCCESS_RESPONSE> resultReader) {
|
||||
protected class RequestPromise<SUCCESS_RESPONSE, RESULT : Any>(private val methodName: String?) : AsyncPromise<RESULT>(), RequestCallback<SUCCESS_RESPONSE> {
|
||||
@Suppress("BASE_WITH_NULLABLE_UPPER_BOUND")
|
||||
override fun onSuccess(response: SUCCESS_RESPONSE?, resultReader: ResultReader<SUCCESS_RESPONSE>?) {
|
||||
try {
|
||||
if (resultReader == null || response == null) {
|
||||
//noinspection unchecked
|
||||
setResult((RESULT)response);
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
setResult(response as RESULT)
|
||||
}
|
||||
else {
|
||||
setResult(methodName == null ? null : resultReader.<RESULT>readResult(methodName, response));
|
||||
if (methodName == null) {
|
||||
setResult(null)
|
||||
}
|
||||
else {
|
||||
setResult(resultReader.readResult(methodName, response))
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (Throwable e) {
|
||||
CommandProcessor.LOG.error(e);
|
||||
setError(e);
|
||||
catch (e: Throwable) {
|
||||
LOG.error(e)
|
||||
setError(e)
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onError(@NotNull Throwable error) {
|
||||
setError(error);
|
||||
override fun onError(error: Throwable) {
|
||||
setError(error)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -13,121 +13,100 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
package org.jetbrains.rpc;
|
||||
package org.jetbrains.rpc
|
||||
|
||||
import com.intellij.util.containers.ConcurrentIntObjectMap;
|
||||
import com.intellij.util.containers.ContainerUtil;
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.concurrency.Promise;
|
||||
import com.intellij.util.containers.ContainerUtil
|
||||
import org.jetbrains.concurrency.Promise
|
||||
import java.io.IOException
|
||||
import java.util.*
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.Arrays;
|
||||
class MessageManager<REQUEST, INCOMING, INCOMING_WITH_SEQ : Any, SUCCESS>(private val handler: MessageManager.Handler<REQUEST, INCOMING, INCOMING_WITH_SEQ, SUCCESS>) : MessageManagerBase() {
|
||||
private val callbackMap = ContainerUtil.createConcurrentIntObjectMap<RequestCallback<SUCCESS>>()
|
||||
|
||||
/**
|
||||
* @param <REQUEST> type of outgoing message
|
||||
* @param <INCOMING> type of incoming message
|
||||
* @param <INCOMING_WITH_SEQ> type of incoming message that is a command (has sequence number)
|
||||
*/
|
||||
public final class MessageManager<REQUEST, INCOMING, INCOMING_WITH_SEQ, SUCCESS> extends MessageManagerBase {
|
||||
private final ConcurrentIntObjectMap<RequestCallback<SUCCESS>> callbackMap = ContainerUtil.createConcurrentIntObjectMap();
|
||||
private final Handler<REQUEST, INCOMING, INCOMING_WITH_SEQ, SUCCESS> handler;
|
||||
interface Handler<OUTGOING, INCOMING, INCOMING_WITH_SEQ : Any, SUCCESS> {
|
||||
fun getUpdatedSequence(message: OUTGOING): Int
|
||||
|
||||
public MessageManager(Handler<REQUEST, INCOMING, INCOMING_WITH_SEQ, SUCCESS> handler) {
|
||||
this.handler = handler;
|
||||
@Throws(IOException::class)
|
||||
fun write(message: OUTGOING): Boolean
|
||||
|
||||
fun readIfHasSequence(incoming: INCOMING): INCOMING_WITH_SEQ?
|
||||
|
||||
fun getSequence(incomingWithSeq: INCOMING_WITH_SEQ): Int
|
||||
|
||||
fun acceptNonSequence(incoming: INCOMING)
|
||||
|
||||
fun call(response: INCOMING_WITH_SEQ, callback: RequestCallback<SUCCESS>)
|
||||
}
|
||||
|
||||
public interface Handler<OUTGOING, INCOMING, INCOMING_WITH_SEQ, SUCCESS> {
|
||||
int getUpdatedSequence(@NotNull OUTGOING message);
|
||||
|
||||
boolean write(@NotNull OUTGOING message) throws IOException;
|
||||
|
||||
INCOMING_WITH_SEQ readIfHasSequence(INCOMING incoming);
|
||||
|
||||
int getSequence(INCOMING_WITH_SEQ incomingWithSeq);
|
||||
|
||||
void acceptNonSequence(INCOMING incoming);
|
||||
|
||||
void call(INCOMING_WITH_SEQ response, RequestCallback<SUCCESS> callback);
|
||||
}
|
||||
|
||||
public void send(@NotNull REQUEST message, @NotNull RequestCallback<SUCCESS> callback) {
|
||||
fun send(message: REQUEST, callback: RequestCallback<SUCCESS>) {
|
||||
if (rejectIfClosed(callback)) {
|
||||
return;
|
||||
return
|
||||
}
|
||||
|
||||
int sequence = handler.getUpdatedSequence(message);
|
||||
callbackMap.put(sequence, callback);
|
||||
|
||||
boolean success;
|
||||
val sequence = handler.getUpdatedSequence(message)
|
||||
callbackMap.put(sequence, callback)
|
||||
|
||||
val success: Boolean
|
||||
try {
|
||||
success = handler.write(message);
|
||||
success = handler.write(message)
|
||||
}
|
||||
catch (Throwable e) {
|
||||
catch (e: Throwable) {
|
||||
try {
|
||||
failedToSend(sequence);
|
||||
failedToSend(sequence)
|
||||
}
|
||||
finally {
|
||||
CommandProcessor.LOG.error("Failed to send", e);
|
||||
LOG.error("Failed to send", e)
|
||||
}
|
||||
return;
|
||||
return
|
||||
}
|
||||
|
||||
if (!success) {
|
||||
failedToSend(sequence);
|
||||
failedToSend(sequence)
|
||||
}
|
||||
}
|
||||
|
||||
private void failedToSend(int sequence) {
|
||||
RequestCallback<SUCCESS> callback = callbackMap.remove(sequence);
|
||||
if (callback != null) {
|
||||
callback.onError(Promise.createError("Failed to send"));
|
||||
}
|
||||
private fun failedToSend(sequence: Int) {
|
||||
callbackMap.remove(sequence)?.onError(Promise.createError("Failed to send"))
|
||||
}
|
||||
|
||||
public void processIncoming(INCOMING incomingParsed) {
|
||||
INCOMING_WITH_SEQ commandResponse = handler.readIfHasSequence(incomingParsed);
|
||||
fun processIncoming(incomingParsed: INCOMING) {
|
||||
val commandResponse = handler.readIfHasSequence(incomingParsed)
|
||||
if (commandResponse == null) {
|
||||
if (closed) {
|
||||
// just ignore
|
||||
CommandProcessor.LOG.info("Connection closed, ignore incoming");
|
||||
LOG.info("Connection closed, ignore incoming")
|
||||
}
|
||||
else {
|
||||
handler.acceptNonSequence(incomingParsed);
|
||||
handler.acceptNonSequence(incomingParsed)
|
||||
}
|
||||
return;
|
||||
return
|
||||
}
|
||||
|
||||
RequestCallback<SUCCESS> callback = getCallbackAndRemove(handler.getSequence(commandResponse));
|
||||
val callback = getCallbackAndRemove(handler.getSequence(commandResponse))
|
||||
if (rejectIfClosed(callback)) {
|
||||
return;
|
||||
return
|
||||
}
|
||||
|
||||
try {
|
||||
handler.call(commandResponse, callback);
|
||||
handler.call(commandResponse, callback)
|
||||
}
|
||||
catch (Throwable e) {
|
||||
callback.onError(e);
|
||||
CommandProcessor.LOG.error("Failed to dispatch response to callback", e);
|
||||
catch (e: Throwable) {
|
||||
callback.onError(e)
|
||||
LOG.error("Failed to dispatch response to callback", e)
|
||||
}
|
||||
}
|
||||
|
||||
public RequestCallback<SUCCESS> getCallbackAndRemove(int id) {
|
||||
RequestCallback<SUCCESS> callback = callbackMap.remove(id);
|
||||
if (callback == null) {
|
||||
throw new IllegalArgumentException("Cannot find callback with id " + id);
|
||||
}
|
||||
return callback;
|
||||
}
|
||||
fun getCallbackAndRemove(id: Int) = callbackMap.remove(id) ?: throw IllegalArgumentException("Cannot find callback with id $id")
|
||||
|
||||
public void cancelWaitingRequests() {
|
||||
fun cancelWaitingRequests() {
|
||||
// we should call them in the order they have been submitted
|
||||
ConcurrentIntObjectMap<RequestCallback<SUCCESS>> map = callbackMap;
|
||||
int[] keys = map.keys();
|
||||
Arrays.sort(keys);
|
||||
for (int key : keys) {
|
||||
RequestCallback<SUCCESS> callback = map.get(key);
|
||||
val map = callbackMap
|
||||
val keys = map.keys()
|
||||
Arrays.sort(keys)
|
||||
for (key in keys) {
|
||||
val callback = map.get(key)
|
||||
if (callback != null) {
|
||||
rejectCallback(callback);
|
||||
MessageManagerBase.rejectCallback(callback)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,24 +1,25 @@
|
||||
package org.jetbrains.rpc;
|
||||
package org.jetbrains.rpc
|
||||
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.concurrency.Promise;
|
||||
import org.jetbrains.concurrency.Promise
|
||||
|
||||
public abstract class MessageManagerBase {
|
||||
protected volatile boolean closed;
|
||||
abstract class MessageManagerBase {
|
||||
@Volatile protected var closed = false
|
||||
|
||||
protected final boolean rejectIfClosed(RequestCallback<?> callback) {
|
||||
protected fun rejectIfClosed(callback: RequestCallback<*>): Boolean {
|
||||
if (closed) {
|
||||
callback.onError(Promise.createError("Connection closed"));
|
||||
return true;
|
||||
callback.onError(Promise.createError("Connection closed"))
|
||||
return true
|
||||
}
|
||||
return false;
|
||||
return false
|
||||
}
|
||||
|
||||
public final void closed() {
|
||||
closed = true;
|
||||
fun closed() {
|
||||
closed = true
|
||||
}
|
||||
|
||||
protected static void rejectCallback(@NotNull RequestCallback<?> callback) {
|
||||
callback.onError(Promise.createError("Connection closed"));
|
||||
companion object {
|
||||
protected fun rejectCallback(callback: RequestCallback<*>) {
|
||||
callback.onError(Promise.createError("Connection closed"))
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,14 +1,12 @@
|
||||
package org.jetbrains.rpc;
|
||||
package org.jetbrains.rpc
|
||||
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.concurrency.Promise;
|
||||
import org.jetbrains.jsonProtocol.Request;
|
||||
import org.jetbrains.concurrency.Promise
|
||||
import org.jetbrains.jsonProtocol.Request
|
||||
|
||||
public interface MessageProcessor {
|
||||
void cancelWaitingRequests();
|
||||
interface MessageProcessor {
|
||||
fun cancelWaitingRequests()
|
||||
|
||||
void closed();
|
||||
fun closed()
|
||||
|
||||
@NotNull
|
||||
<T> Promise<T> send(@NotNull Request<T> message);
|
||||
fun <RESULT : Any> send(message: Request<RESULT>): Promise<RESULT>
|
||||
}
|
||||
@@ -1,26 +1,21 @@
|
||||
package org.jetbrains.rpc;
|
||||
package org.jetbrains.rpc
|
||||
|
||||
import com.intellij.openapi.vfs.CharsetToolkit;
|
||||
import com.intellij.util.BooleanFunction;
|
||||
import io.netty.buffer.ByteBuf;
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.jsonProtocol.Request;
|
||||
import com.intellij.openapi.vfs.CharsetToolkit
|
||||
import com.intellij.util.BooleanFunction
|
||||
import io.netty.buffer.ByteBuf
|
||||
import org.jetbrains.jsonProtocol.Request
|
||||
|
||||
import static org.jetbrains.rpc.CommandProcessor.LOG;
|
||||
|
||||
public abstract class MessageWriter implements BooleanFunction<Request> {
|
||||
@Override
|
||||
public boolean fun(@NotNull Request message) {
|
||||
ByteBuf content = message.getBuffer();
|
||||
if (isDebugLoggingEnabled()) {
|
||||
LOG.debug("OUT: " + content.toString(CharsetToolkit.UTF8_CHARSET));
|
||||
abstract class MessageWriter : BooleanFunction<Request<Any>> {
|
||||
override fun `fun`(message: Request<Any>): Boolean {
|
||||
val content = message.buffer
|
||||
if (isDebugLoggingEnabled) {
|
||||
LOG.debug("OUT: ${content.toString(CharsetToolkit.UTF8_CHARSET)}")
|
||||
}
|
||||
return write(content);
|
||||
return write(content)
|
||||
}
|
||||
|
||||
protected boolean isDebugLoggingEnabled() {
|
||||
return LOG.isDebugEnabled();
|
||||
}
|
||||
protected open val isDebugLoggingEnabled: Boolean
|
||||
get() = LOG.isDebugEnabled
|
||||
|
||||
protected abstract boolean write(@NotNull ByteBuf content);
|
||||
protected abstract fun write(content: ByteBuf): Boolean
|
||||
}
|
||||
@@ -1,10 +1,8 @@
|
||||
package org.jetbrains.rpc;
|
||||
package org.jetbrains.rpc
|
||||
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.annotations.Nullable;
|
||||
@Suppress("BASE_WITH_NULLABLE_UPPER_BOUND")
|
||||
interface RequestCallback<SUCCESS_RESPONSE> {
|
||||
fun onSuccess(response: SUCCESS_RESPONSE?, resultReader: ResultReader<SUCCESS_RESPONSE>?)
|
||||
|
||||
public interface RequestCallback<SUCCESS_RESPONSE> {
|
||||
void onSuccess(@Nullable SUCCESS_RESPONSE successResponse, @Nullable ResultReader<SUCCESS_RESPONSE> resultReader);
|
||||
|
||||
void onError(@NotNull Throwable error);
|
||||
fun onError(error: Throwable)
|
||||
}
|
||||
@@ -1,7 +1,5 @@
|
||||
package org.jetbrains.rpc;
|
||||
package org.jetbrains.rpc
|
||||
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
|
||||
public interface ResultReader<RESPONSE> {
|
||||
<RESULT> RESULT readResult(@NotNull String readMethodName, @NotNull RESPONSE successResponse);
|
||||
interface ResultReader<RESPONSE> {
|
||||
fun <RESULT> readResult(readMethodName: String, successResponse: RESPONSE): RESULT
|
||||
}
|
||||
+15
-15
@@ -13,21 +13,21 @@
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
package org.jetbrains.debugger;
|
||||
package org.jetbrains.debugger
|
||||
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.concurrency.AsyncFunction;
|
||||
import org.jetbrains.concurrency.Obsolescent;
|
||||
import org.jetbrains.concurrency.AsyncFunction
|
||||
import org.jetbrains.concurrency.Obsolescent
|
||||
import org.jetbrains.concurrency.Promise
|
||||
|
||||
public abstract class ValueNodeAsyncFunction<PARAM, RESULT> implements AsyncFunction<PARAM, RESULT>, Obsolescent {
|
||||
private final Obsolescent node;
|
||||
|
||||
protected ValueNodeAsyncFunction(@NotNull Obsolescent node) {
|
||||
this.node = node;
|
||||
}
|
||||
|
||||
@Override
|
||||
public final boolean isObsolete() {
|
||||
return node.isObsolete();
|
||||
}
|
||||
abstract class ValueNodeAsyncFunction<PARAM, RESULT> protected constructor(private val node: Obsolescent) : AsyncFunction<PARAM, RESULT>, Obsolescent {
|
||||
override fun isObsolete() = node.isObsolete
|
||||
}
|
||||
|
||||
inline fun <T, SUB_RESULT> Promise<T>.thenAsync(node: Obsolescent, crossinline handler: (T) -> Promise<SUB_RESULT>) = then(object : ValueNodeAsyncFunction<T, SUB_RESULT>(node) {
|
||||
override fun `fun`(param: T) = handler(param)
|
||||
})
|
||||
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
inline fun <T> Promise<T>.thenAsyncVoid(node: Obsolescent, crossinline handler: (T) -> Promise<*>) = then(object : ValueNodeAsyncFunction<T, Any?>(node) {
|
||||
override fun `fun`(param: T) = handler(param) as Promise<Any?>
|
||||
})
|
||||
Reference in New Issue
Block a user