rpc client — use promises

efficient MessageDecoder — don't create String instance
simplify Decoder implementations — base Decoder now provides robust readContent API
MessageDecoder extends Decoder and provides robust readChars API
This commit is contained in:
Vladimir Krivosheev
2015-02-15 11:51:54 +01:00
parent 2d444c02da
commit 285fdcfc82
72 changed files with 1997 additions and 861 deletions
+2 -2
View File
@@ -1,11 +1,11 @@
<component name="libraryTable">
<library name="Netty">
<CLASSES>
<root url="jar://$PROJECT_DIR$/lib/netty-all-4.1.0.Beta3.jar!/" />
<root url="jar://$PROJECT_DIR$/lib/netty-all-4.1.0.Beta4.jar!/" />
</CLASSES>
<JAVADOC />
<SOURCES>
<root url="jar://$PROJECT_DIR$/lib/src/netty-all-4.1.0.Beta3-sources.jar!/" />
<root url="jar://$PROJECT_DIR$/lib/src/netty-all-4.1.0.Beta4-sources.jar!/" />
</SOURCES>
</library>
</component>
+5
View File
@@ -0,0 +1,5 @@
<root>
<item name='io.netty.buffer.ByteBuf java.lang.String toString(int, int, java.nio.charset.Charset)'>
<annotation name='org.jetbrains.annotations.NotNull'/>
</item>
</root>
+3
View File
@@ -1,4 +1,7 @@
<root>
<item name='io.netty.channel.ChannelInboundHandlerAdapter void channelInactive(io.netty.channel.ChannelHandlerContext) 0'>
<annotation name='org.jetbrains.annotations.NotNull'/>
</item>
<item name='io.netty.channel.ChannelInitializer void channelRegistered(io.netty.channel.ChannelHandlerContext) 0'>
<annotation name='org.jetbrains.annotations.NotNull'/>
</item>
Binary file not shown.
Binary file not shown.
+1 -1
View File
@@ -48,7 +48,7 @@ microba.jar
miglayout-swing.jar
nanoxml-2.2.3.jar
nekohtml-1.9.14.jar
netty-all-4.1.0.Beta3.jar
netty-all-4.1.0.Beta4.jar
oromatcher.jar
picocontainer.jar
protobuf-2.5.0.jar
@@ -17,12 +17,13 @@ package org.jetbrains.ide;
import com.intellij.openapi.extensions.ExtensionPointName;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelHandlerContext;
import org.jetbrains.annotations.NotNull;
import java.util.UUID;
public abstract class BinaryRequestHandler {
public static final ExtensionPointName<BinaryRequestHandler> EP_NAME = ExtensionPointName.create("com.intellij.binaryRequestHandler");
public static final ExtensionPointName<BinaryRequestHandler> EP_NAME = ExtensionPointName.create("org.jetbrains.binaryRequestHandler");
@NotNull
/**
@@ -31,5 +32,5 @@ public abstract class BinaryRequestHandler {
public abstract UUID getId();
@NotNull
public abstract ChannelHandler getInboundHandler();
public abstract ChannelHandler getInboundHandler(@NotNull ChannelHandlerContext context);
}
@@ -25,12 +25,12 @@ gulp.task("compile", function () {
.pipe(ts(tsProject));
tsResult.js.pipe(concat(outFile))
.pipe(uglify({
output: {
beautify: true,
indent_level: 2
}
}))
//.pipe(uglify({
// output: {
// beautify: true,
// indent_level: 2
// }
// }))
.pipe(sourcemaps.write('.', {includeContent: false, sourceRoot: path.resolve('testData')}))
.pipe(gulp.dest(outDir))
});
@@ -1,7 +1,11 @@
{
"name": "ij-rpc-client",
"version": 0.1,
"version": "0.0.1",
"description": "IntelliJ Platform RPC client",
"repository": {
"type": "git",
"url": "https://github.com/JetBrains/intellij-community.git"
},
"devDependencies": {
"gulp": "^3.8.11",
"gulp-concat": "^2.4.3",
@@ -1,6 +1,7 @@
/// <reference path="../typings/node/node.d.ts" />
"use strict"
import net = require("net")
import rpc = require("./rpc")
export class RpcClient {
connect(port:Number = 63342) {
@@ -8,92 +9,43 @@ export class RpcClient {
console.log('Connected to IJ RPC server localhost:' + port)
});
var jsonRpc = new JsonRpc(new Socket(socket))
var jsonRpc = new rpc.JsonRpc(new SocketTransport(socket))
var decoder:MessageDecoder = new MessageDecoder(jsonRpc.messageReceived)
socket.on('data', decoder.messageReceived)
}
}
enum State {LENGTH, CONTENT}
const enum State {LENGTH, CONTENT}
class Socket {
private uint32Buffer = new Buffer(4)
class SocketTransport implements rpc.Transport {
private headerBuffer = new Buffer(4)
constructor(private socket:net.Socket) {
}
send(message:string) {
var encoded = new Buffer(message)
this.uint32Buffer.writeUInt32BE(Buffer.byteLength(message), 0)
this.socket.write(this.uint32Buffer)
this.socket.write(message, 'utf-8')
}
}
class JsonRpc {
private messageIdCounter = 0
private callbacks:Map<number, (result?:any)=>void> = new Map<number, (result:any)=>void>()
private domains:Map<string, any> = new Map<string, any>();
constructor(private socket:Socket) {
send(id:number, domain:string, command:string, params:any[] = null):void {
var encodedParams = JSON.stringify(params)
var header = (id == -1 ? '' : (id + ', ')) + '"' + domain + '", "' + command + '"';
this.headerBuffer.writeUInt32BE(Buffer.byteLength(encodedParams) + header.length, 0)
this.socket.write(this.headerBuffer)
this.socket.write(encodedParams, 'utf-8')
}
messageReceived(message:Array<any>) {
if (message.length === 1 || (message.length === 2 && !(typeof message[1] === 'string'))) {
var f = this.callbacks.get(message[0])
var singletonArray = safeGet(message, 1)
if (singletonArray == null) {
f()
}
else {
f(singletonArray[0])
}
}
else {
var id:number
var offset:number
if (typeof message[0] === 'string') {
id = -1
offset = 0
}
else {
id = message[0]
offset = 1
}
var args = safeGet(message, offset + 2)
var errorCallback:(message:String)=>void
if (id === -1) {
errorCallback = null
}
else {
var resultCallback = (result:any) => this.socket.send("[" + id + ", \"r\", " + JSON.stringify(result) + "]")
errorCallback = (error:any) => this.socket.send("[$id, \"e\", " + JSON.stringify(error) + "]")
if (args == null) {
args = [resultCallback, errorCallback]
}
else {
var regularArgs = args
args = regularArgs.concat(errorCallback, resultCallback)
}
}
try {
var o = this.domains.get(message[offset]);
o[message[offset + 1]].apply(o, args);
}
catch (e) {
console.error(e)
if (errorCallback != null) {
errorCallback(e)
}
}
}
sendResult(id:number, result:any):void {
this.sendResultOrError(id, result, false)
}
}
function safeGet(a:Array<any>, index:number):Array<any> {
return index < a.length ? a[index] : null
sendError(id:number, error:any):void {
this.sendResultOrError(id, error, true)
}
private sendResultOrError(id:number, result:any, isError:boolean):void {
var encodedResult = JSON.stringify(result)
var header = id + ', "' + (isError ? 'e': 'r') + '"';
this.headerBuffer.writeUInt32BE(Buffer.byteLength(encodedResult) + header.length, 0)
this.socket.write(this.headerBuffer)
this.socket.write(encodedResult, 'utf-8')
}
}
class MessageDecoder {
@@ -0,0 +1,94 @@
/// <reference path="../typings/node/node.d.ts" />
/// <reference path="../typings/bluebird/bluebird.d.ts" />
"use strict"
import Promise = require("bluebird")
class PromiseCallback {
constructor(public resolve:(value?:any) => void, public reject:(error?:any) => void) {
}
}
export interface Transport {
send(id:number, domain:string, command:string, params:any[]):void
sendResult(id:number, result:any):void
sendError(id:number, error:any):void
}
export class JsonRpc {
private messageIdCounter = 0
private callbacks:Map<number, PromiseCallback> = new Map<number, PromiseCallback>()
private domains:Map<string, any> = new Map<string, any>()
constructor(private transport:Transport) {
}
public call<T>(domain:string, command:string, ...params: any[]):Promise<T> {
return new Promise((resolve:(value:T) => void, reject:(error?:any) => void) => {
var id = this.messageIdCounter++;
this.callbacks.set(id, new PromiseCallback(resolve, reject))
this.transport.send(id, domain, command, params)
})
}
messageReceived(message:Array<any>) {
if (message.length === 1 || (message.length === 2 && !(typeof message[1] === 'string'))) {
var promiseCallback = this.callbacks.get(message[0])
var singletonArray = safeGet(message, 1)
if (singletonArray == null) {
promiseCallback.resolve()
}
else {
promiseCallback.resolve(singletonArray[0])
}
}
else {
var id:number
var offset:number
if (typeof message[0] === 'string') {
id = -1
offset = 0
}
else {
id = message[0]
offset = 1
}
var onRejected = id === -1 ? null : (error:any) => this.transport.sendError(id, error)
try {
var object = this.domains.get(message[offset])
var method = object[message[offset + 1]]
var result:any
var args = safeGet(message, offset + 2)
if (args === null) {
result = method.call(object)
}
else {
result = method.apply(object, args)
}
if (id !== -1) {
var onFulfilled = (result:any) => this.transport.sendResult(id, result)
if (result instanceof Promise) {
(<Promise<any>>result).done(onFulfilled, onRejected)
}
else {
onFulfilled(result)
}
}
}
catch (e) {
console.error(e)
if (onRejected != null) {
onRejected(e)
}
}
}
}
}
function safeGet(a:any[], index:number):Array<any> {
return index < a.length ? a[index] : null
}
@@ -7,6 +7,9 @@
"installed": {
"node/node.d.ts": {
"commit": "35fffaa44bff5392741b3022d805fe3563023a3d"
},
"bluebird/bluebird.d.ts": {
"commit": "cf7c97b2a68a385c98c75fb6edd81083c97c983c"
}
}
}
@@ -0,0 +1,710 @@
// Type definitions for bluebird 2.0.0
// Project: https://github.com/petkaantonov/bluebird
// Definitions by: Bart van der Schoor <https://github.com/Bartvds>
// Definitions: https://github.com/borisyankov/DefinitelyTyped
// ES6 model with generics overload was sourced and trans-multiplied from es6-promises.d.ts
// By: Campredon <https://github.com/fdecampredon/>
// Warning: recommended to use `tsc > v0.9.7` (critical bugs in earlier generic code):
// - https://github.com/borisyankov/DefinitelyTyped/issues/1563
// Note: replicate changes to all overloads in both definition and test file
// Note: keep both static and instance members inline (so similar)
// TODO fix remaining TODO annotations in both definition and test
// TODO verify support to have no return statement in handlers to get a Promise<void> (more overloads?)
declare class Promise<R> implements Promise.Thenable<R>, Promise.Inspection<R> {
/**
* Create a new promise. The passed in function will receive functions `resolve` and `reject` as its arguments which can be called to seal the fate of the created promise.
*/
constructor(callback: (resolve: (thenable: Promise.Thenable<R>) => void, reject: (error: any) => void) => void);
constructor(callback: (resolve: (result: R) => void, reject: (error: any) => void) => void);
/**
* Promises/A+ `.then()` with progress handler. Returns a new promise chained from this promise. The new promise will be rejected or resolved dedefer on the passed `fulfilledHandler`, `rejectedHandler` and the state of this promise.
*/
then<U>(onFulfill: (value: R) => Promise.Thenable<U>, onReject: (error: any) => Promise.Thenable<U>, onProgress?: (note: any) => any): Promise<U>;
then<U>(onFulfill: (value: R) => Promise.Thenable<U>, onReject?: (error: any) => U, onProgress?: (note: any) => any): Promise<U>;
then<U>(onFulfill: (value: R) => U, onReject: (error: any) => Promise.Thenable<U>, onProgress?: (note: any) => any): Promise<U>;
then<U>(onFulfill?: (value: R) => U, onReject?: (error: any) => U, onProgress?: (note: any) => any): Promise<U>;
/**
* This is a catch-all exception handler, shortcut for calling `.then(null, handler)` on this promise. Any exception happening in a `.then`-chain will propagate to nearest `.catch` handler.
*
* Alias `.caught();` for compatibility with earlier ECMAScript version.
*/
catch<U>(onReject?: (error: any) => Promise.Thenable<U>): Promise<U>;
caught<U>(onReject?: (error: any) => Promise.Thenable<U>): Promise<U>;
catch<U>(onReject?: (error: any) => U): Promise<U>;
caught<U>(onReject?: (error: any) => U): Promise<U>;
/**
* This extends `.catch` to work more like catch-clauses in languages like Java or C#. Instead of manually checking `instanceof` or `.name === "SomeError"`, you may specify a number of error constructors which are eligible for this catch handler. The catch handler that is first met that has eligible constructors specified, is the one that will be called.
*
* This method also supports predicate-based filters. If you pass a predicate function instead of an error constructor, the predicate will receive the error as an argument. The return result of the predicate will be used determine whether the error handler should be called.
*
* Alias `.caught();` for compatibility with earlier ECMAScript version.
*/
catch<U>(predicate: (error: any) => boolean, onReject: (error: any) => Promise.Thenable<U>): Promise<U>;
caught<U>(predicate: (error: any) => boolean, onReject: (error: any) => Promise.Thenable<U>): Promise<U>;
catch<U>(predicate: (error: any) => boolean, onReject: (error: any) => U): Promise<U>;
caught<U>(predicate: (error: any) => boolean, onReject: (error: any) => U): Promise<U>;
catch<U>(ErrorClass: Function, onReject: (error: any) => Promise.Thenable<U>): Promise<U>;
caught<U>(ErrorClass: Function, onReject: (error: any) => Promise.Thenable<U>): Promise<U>;
catch<U>(ErrorClass: Function, onReject: (error: any) => U): Promise<U>;
caught<U>(ErrorClass: Function, onReject: (error: any) => U): Promise<U>;
/**
* Like `.catch` but instead of catching all types of exceptions, it only catches those that don't originate from thrown errors but rather from explicit rejections.
*/
error<U>(onReject: (reason: any) => Promise.Thenable<U>): Promise<U>;
error<U>(onReject: (reason: any) => U): Promise<U>;
/**
* Pass a handler that will be called regardless of this promise's fate. Returns a new promise chained from this promise. There are special semantics for `.finally()` in that the final value cannot be modified from the handler.
*
* Alias `.lastly();` for compatibility with earlier ECMAScript version.
*/
finally<U>(handler: () => Promise.Thenable<U>): Promise<R>;
finally<U>(handler: () => U): Promise<R>;
lastly<U>(handler: () => Promise.Thenable<U>): Promise<R>;
lastly<U>(handler: () => U): Promise<R>;
/**
* Create a promise that follows this promise, but is bound to the given `thisArg` value. A bound promise will call its handlers with the bound value set to `this`. Additionally promises derived from a bound promise will also be bound promises with the same `thisArg` binding as the original promise.
*/
bind(thisArg: any): Promise<R>;
/**
* Like `.then()`, but any unhandled rejection that ends up here will be thrown as an error.
*/
done<U>(onFulfilled: (value: R) => Promise.Thenable<U>, onRejected: (error: any) => Promise.Thenable<U>, onProgress?: (note: any) => any): void;
done<U>(onFulfilled: (value: R) => Promise.Thenable<U>, onRejected?: (error: any) => U, onProgress?: (note: any) => any): void;
done<U>(onFulfilled: (value: R) => U, onRejected: (error: any) => Promise.Thenable<U>, onProgress?: (note: any) => any): void;
done<U>(onFulfilled?: (value: R) => U, onRejected?: (error: any) => U, onProgress?: (note: any) => any): void;
/**
* Like `.finally()`, but not called for rejections.
*/
tap<U>(onFulFill: (value: R) => Promise.Thenable<U>): Promise<R>;
tap<U>(onFulfill: (value: R) => U): Promise<R>;
/**
* Shorthand for `.then(null, null, handler);`. Attach a progress handler that will be called if this promise is progressed. Returns a new promise chained from this promise.
*/
progressed(handler: (note: any) => any): Promise<R>;
/**
* Same as calling `Promise.delay(this, ms)`. With the exception that if this promise is bound to a value, the returned promise is bound to that value too.
*/
delay(ms: number): Promise<R>;
/**
* Returns a promise that will be fulfilled with this promise's fulfillment value or rejection reason. However, if this promise is not fulfilled or rejected within `ms` milliseconds, the returned promise is rejected with a `Promise.TimeoutError` instance.
*
* You may specify a custom error message with the `message` parameter.
*/
timeout(ms: number, message?: string): Promise<R>;
/**
* Register a node-style callback on this promise. When this promise is is either fulfilled or rejected, the node callback will be called back with the node.js convention where error reason is the first argument and success value is the second argument. The error argument will be `null` in case of success.
* Returns back this promise instead of creating a new one. If the `callback` argument is not a function, this method does not do anything.
*/
nodeify(callback: (err: any, value?: R) => void): Promise<R>;
nodeify(...sink: any[]): void;
/**
* Marks this promise as cancellable. Promises by default are not cancellable after v0.11 and must be marked as such for `.cancel()` to have any effect. Marking a promise as cancellable is infectious and you don't need to remark any descendant promise.
*/
cancellable(): Promise<R>;
/**
* Cancel this promise. The cancellation will propagate to farthest cancellable ancestor promise which is still pending.
*
* That ancestor will then be rejected with a `CancellationError` (get a reference from `Promise.CancellationError`) object as the rejection reason.
*
* In a promise rejection handler you may check for a cancellation by seeing if the reason object has `.name === "Cancel"`.
*
* Promises are by default not cancellable. Use `.cancellable()` to mark a promise as cancellable.
*/
// TODO what to do with this?
cancel<U>(): Promise<U>;
/**
* Like `.then()`, but cancellation of the the returned promise or any of its descendant will not propagate cancellation to this promise or this promise's ancestors.
*/
fork<U>(onFulfilled: (value: R) => Promise.Thenable<U>, onRejected: (error: any) => Promise.Thenable<U>, onProgress?: (note: any) => any): Promise<U>;
fork<U>(onFulfilled: (value: R) => Promise.Thenable<U>, onRejected?: (error: any) => U, onProgress?: (note: any) => any): Promise<U>;
fork<U>(onFulfilled: (value: R) => U, onRejected: (error: any) => Promise.Thenable<U>, onProgress?: (note: any) => any): Promise<U>;
fork<U>(onFulfilled?: (value: R) => U, onRejected?: (error: any) => U, onProgress?: (note: any) => any): Promise<U>;
/**
* Create an uncancellable promise based on this promise.
*/
uncancellable(): Promise<R>;
/**
* See if this promise can be cancelled.
*/
isCancellable(): boolean;
/**
* See if this `promise` has been fulfilled.
*/
isFulfilled(): boolean;
/**
* See if this `promise` has been rejected.
*/
isRejected(): boolean;
/**
* See if this `promise` is still defer.
*/
isPending(): boolean;
/**
* See if this `promise` is resolved -> either fulfilled or rejected.
*/
isResolved(): boolean;
/**
* Get the fulfillment value of the underlying promise. Throws if the promise isn't fulfilled yet.
*
* throws `TypeError`
*/
value(): R;
/**
* Get the rejection reason for the underlying promise. Throws if the promise isn't rejected yet.
*
* throws `TypeError`
*/
reason(): any;
/**
* Synchronously inspect the state of this `promise`. The `PromiseInspection` will represent the state of the promise as snapshotted at the time of calling `.inspect()`.
*/
inspect(): Promise.Inspection<R>;
/**
* This is a convenience method for doing:
*
* <code>
* promise.then(function(obj){
* return obj[propertyName].call(obj, arg...);
* });
* </code>
*/
call(propertyName: string, ...args: any[]): Promise<any>;
/**
* This is a convenience method for doing:
*
* <code>
* promise.then(function(obj){
* return obj[propertyName];
* });
* </code>
*/
// TODO find way to fix get()
// get<U>(propertyName: string): Promise<U>;
/**
* Convenience method for:
*
* <code>
* .then(function() {
* return value;
* });
* </code>
*
* in the case where `value` doesn't change its value. That means `value` is bound at the time of calling `.return()`
*
* Alias `.thenReturn();` for compatibility with earlier ECMAScript version.
*/
return(): Promise<any>;
thenReturn(): Promise<any>;
return<U>(value: U): Promise<U>;
thenReturn<U>(value: U): Promise<U>;
/**
* Convenience method for:
*
* <code>
* .then(function() {
* throw reason;
* });
* </code>
* Same limitations apply as with `.return()`.
*
* Alias `.thenThrow();` for compatibility with earlier ECMAScript version.
*/
throw(reason: Error): Promise<R>;
thenThrow(reason: Error): Promise<R>;
/**
* Convert to String.
*/
toString(): string;
/**
* This is implicitly called by `JSON.stringify` when serializing the object. Returns a serialized representation of the `Promise`.
*/
toJSON(): Object;
/**
* Like calling `.then`, but the fulfillment value or rejection reason is assumed to be an array, which is flattened to the formal parameters of the handlers.
*/
// TODO how to model instance.spread()? like Q?
spread<U>(onFulfill: Function, onReject?: (reason: any) => Promise.Thenable<U>): Promise<U>;
spread<U>(onFulfill: Function, onReject?: (reason: any) => U): Promise<U>;
/*
// TODO or something like this?
spread<U, W>(onFulfill: (...values: W[]) => Promise.Thenable<U>, onReject?: (reason: any) => Promise.Thenable<U>): Promise<U>;
spread<U, W>(onFulfill: (...values: W[]) => Promise.Thenable<U>, onReject?: (reason: any) => U): Promise<U>;
spread<U, W>(onFulfill: (...values: W[]) => U, onReject?: (reason: any) => Promise.Thenable<U>): Promise<U>;
spread<U, W>(onFulfill: (...values: W[]) => U, onReject?: (reason: any) => U): Promise<U>;
*/
/**
* Same as calling `Promise.all(thisPromise)`. With the exception that if this promise is bound to a value, the returned promise is bound to that value too.
*/
// TODO type inference from array-resolving promise?
all<U>(): Promise<U[]>;
/**
* Same as calling `Promise.props(thisPromise)`. With the exception that if this promise is bound to a value, the returned promise is bound to that value too.
*/
// TODO how to model instance.props()?
props(): Promise<Object>;
/**
* Same as calling `Promise.settle(thisPromise)`. With the exception that if this promise is bound to a value, the returned promise is bound to that value too.
*/
// TODO type inference from array-resolving promise?
settle<U>(): Promise<Promise.Inspection<U>[]>;
/**
* Same as calling `Promise.any(thisPromise)`. With the exception that if this promise is bound to a value, the returned promise is bound to that value too.
*/
// TODO type inference from array-resolving promise?
any<U>(): Promise<U>;
/**
* Same as calling `Promise.some(thisPromise)`. With the exception that if this promise is bound to a value, the returned promise is bound to that value too.
*/
// TODO type inference from array-resolving promise?
some<U>(count: number): Promise<U[]>;
/**
* Same as calling `Promise.race(thisPromise, count)`. With the exception that if this promise is bound to a value, the returned promise is bound to that value too.
*/
// TODO type inference from array-resolving promise?
race<U>(): Promise<U>;
/**
* Same as calling `Promise.map(thisPromise, mapper)`. With the exception that if this promise is bound to a value, the returned promise is bound to that value too.
*/
// TODO type inference from array-resolving promise?
map<Q, U>(mapper: (item: Q, index: number, arrayLength: number) => Promise.Thenable<U>): Promise<U[]>;
map<Q, U>(mapper: (item: Q, index: number, arrayLength: number) => U): Promise<U[]>;
/**
* Same as calling `Promise.reduce(thisPromise, Function reducer, initialValue)`. With the exception that if this promise is bound to a value, the returned promise is bound to that value too.
*/
// TODO type inference from array-resolving promise?
reduce<Q, U>(reducer: (memo: U, item: Q, index: number, arrayLength: number) => Promise.Thenable<U>, initialValue?: U): Promise<U>;
reduce<Q, U>(reducer: (memo: U, item: Q, index: number, arrayLength: number) => U, initialValue?: U): Promise<U>;
/**
* Same as calling ``Promise.filter(thisPromise, filterer)``. With the exception that if this promise is bound to a value, the returned promise is bound to that value too.
*/
// TODO type inference from array-resolving promise?
filter<U>(filterer: (item: U, index: number, arrayLength: number) => Promise.Thenable<boolean>): Promise<U[]>;
filter<U>(filterer: (item: U, index: number, arrayLength: number) => boolean): Promise<U[]>;
/**
* Start the chain of promises with `Promise.try`. Any synchronous exceptions will be turned into rejections on the returned promise.
*
* Note about second argument: if it's specifically a true array, its values become respective arguments for the function call. Otherwise it is passed as is as the first argument for the function call.
*
* Alias for `attempt();` for compatibility with earlier ECMAScript version.
*/
static try<R>(fn: () => Promise.Thenable<R>, args?: any[], ctx?: any): Promise<R>;
static try<R>(fn: () => R, args?: any[], ctx?: any): Promise<R>;
static attempt<R>(fn: () => Promise.Thenable<R>, args?: any[], ctx?: any): Promise<R>;
static attempt<R>(fn: () => R, args?: any[], ctx?: any): Promise<R>;
/**
* Returns a new function that wraps the given function `fn`. The new function will always return a promise that is fulfilled with the original functions return values or rejected with thrown exceptions from the original function.
* This method is convenient when a function can sometimes return synchronously or throw synchronously.
*/
static method(fn: Function): Function;
/**
* Create a promise that is resolved with the given `value`. If `value` is a thenable or promise, the returned promise will assume its state.
*/
static resolve(): Promise<void>;
static resolve<R>(value: Promise.Thenable<R>): Promise<R>;
static resolve<R>(value: R): Promise<R>;
/**
* Create a promise that is rejected with the given `reason`.
*/
static reject(reason: any): Promise<any>;
static reject<R>(reason: any): Promise<R>;
/**
* Create a promise with undecided fate and return a `PromiseResolver` to control it. See resolution?: Promise(#promise-resolution).
*/
static defer<R>(): Promise.Resolver<R>;
/**
* Cast the given `value` to a trusted promise. If `value` is already a trusted `Promise`, it is returned as is. If `value` is not a thenable, a fulfilled is: Promise returned with `value` as its fulfillment value. If `value` is a thenable (Promise-like object, like those returned by jQuery's `$.ajax`), returns a trusted that: Promise assimilates the state of the thenable.
*/
static cast<R>(value: Promise.Thenable<R>): Promise<R>;
static cast<R>(value: R): Promise<R>;
/**
* Sugar for `Promise.resolve(undefined).bind(thisArg);`. See `.bind()`.
*/
static bind(thisArg: any): Promise<void>;
/**
* See if `value` is a trusted Promise.
*/
static is(value: any): boolean;
/**
* Call this right after the library is loaded to enabled long stack traces. Long stack traces cannot be disabled after being enabled, and cannot be enabled after promises have alread been created. Long stack traces imply a substantial performance penalty, around 4-5x for throughput and 0.5x for latency.
*/
static longStackTraces(): void;
/**
* Returns a promise that will be fulfilled with `value` (or `undefined`) after given `ms` milliseconds. If `value` is a promise, the delay will start counting down when it is fulfilled and the returned promise will be fulfilled with the fulfillment value of the `value` promise.
*/
// TODO enable more overloads
static delay<R>(value: Promise.Thenable<R>, ms: number): Promise<R>;
static delay<R>(value: R, ms: number): Promise<R>;
static delay(ms: number): Promise<void>;
/**
* Returns a function that will wrap the given `nodeFunction`. Instead of taking a callback, the returned function will return a promise whose fate is decided by the callback behavior of the given node function. The node function should conform to node.js convention of accepting a callback as last argument and calling that callback with error as the first argument and success value on the second argument.
*
* If the `nodeFunction` calls its callback with multiple success values, the fulfillment value will be an array of them.
*
* If you pass a `receiver`, the `nodeFunction` will be called as a method on the `receiver`.
*/
// TODO how to model promisify?
static promisify(nodeFunction: Function, receiver?: any): Function;
/**
* Promisifies the entire object by going through the object's properties and creating an async equivalent of each function on the object and its prototype chain. The promisified method name will be the original method name postfixed with `Async`. Returns the input object.
*
* Note that the original methods on the object are not overwritten but new methods are created with the `Async`-postfix. For example, if you `promisifyAll()` the node.js `fs` object use `fs.statAsync()` to call the promisified `stat` method.
*/
// TODO how to model promisifyAll?
static promisifyAll(target: Object): Object;
/**
* Returns a function that can use `yield` to run asynchronous code synchronously. This feature requires the support of generators which are drafted in the next version of the language. Node version greater than `0.11.2` is required and needs to be executed with the `--harmony-generators` (or `--harmony`) command-line switch.
*/
// TODO fix coroutine GeneratorFunction
static coroutine<R>(generatorFunction: Function): Function;
/**
* Spawn a coroutine which may yield promises to run asynchronous code synchronously. This feature requires the support of generators which are drafted in the next version of the language. Node version greater than `0.11.2` is required and needs to be executed with the `--harmony-generators` (or `--harmony`) command-line switch.
*/
// TODO fix spawn GeneratorFunction
static spawn<R>(generatorFunction: Function): Promise<R>;
/**
* This is relevant to browser environments with no module loader.
*
* Release control of the `Promise` namespace to whatever it was before this library was loaded. Returns a reference to the library namespace so you can attach it to something else.
*/
static noConflict(): typeof Promise;
/**
* Add `handler` as the handler to call when there is a possibly unhandled rejection. The default handler logs the error stack to stderr or `console.error` in browsers.
*
* Passing no value or a non-function will have the effect of removing any kind of handling for possibly unhandled rejections.
*/
static onPossiblyUnhandledRejection(handler: (reason: any) => any): void;
/**
* Given an array, or a promise of an array, which contains promises (or a mix of promises and values) return a promise that is fulfilled when all the items in the array are fulfilled. The promise's fulfillment value is an array with fulfillment values at respective positions to the original array. If any promise in the array rejects, the returned promise is rejected with the rejection reason.
*/
// TODO enable more overloads
// promise of array with promises of value
static all<R>(values: Promise.Thenable<Promise.Thenable<R>[]>): Promise<R[]>;
// promise of array with values
static all<R>(values: Promise.Thenable<R[]>): Promise<R[]>;
// array with promises of value
static all<R>(values: Promise.Thenable<R>[]): Promise<R[]>;
// array with values
static all<R>(values: R[]): Promise<R[]>;
/**
* Like ``Promise.all`` but for object properties instead of array items. Returns a promise that is fulfilled when all the properties of the object are fulfilled. The promise's fulfillment value is an object with fulfillment values at respective keys to the original object. If any promise in the object rejects, the returned promise is rejected with the rejection reason.
*
* If `object` is a trusted `Promise`, then it will be treated as a promise for object rather than for its properties. All other objects are treated for their properties as is returned by `Object.keys` - the object's own enumerable properties.
*
* *The original object is not modified.*
*/
// TODO verify this is correct
// trusted promise for object
static props(object: Promise<Object>): Promise<Object>;
// object
static props(object: Object): Promise<Object>;
/**
* Given an array, or a promise of an array, which contains promises (or a mix of promises and values) return a promise that is fulfilled when all the items in the array are either fulfilled or rejected. The fulfillment value is an array of ``PromiseInspection`` instances at respective positions in relation to the input array.
*
* *original: The array is not modified. The input array sparsity is retained in the resulting array.*
*/
// promise of array with promises of value
static settle<R>(values: Promise.Thenable<Promise.Thenable<R>[]>): Promise<Promise.Inspection<R>[]>;
// promise of array with values
static settle<R>(values: Promise.Thenable<R[]>): Promise<Promise.Inspection<R>[]>;
// array with promises of value
static settle<R>(values: Promise.Thenable<R>[]): Promise<Promise.Inspection<R>[]>;
// array with values
static settle<R>(values: R[]): Promise<Promise.Inspection<R>[]>;
/**
* Like `Promise.some()`, with 1 as `count`. However, if the promise fulfills, the fulfillment value is not an array of 1 but the value directly.
*/
// promise of array with promises of value
static any<R>(values: Promise.Thenable<Promise.Thenable<R>[]>): Promise<R>;
// promise of array with values
static any<R>(values: Promise.Thenable<R[]>): Promise<R>;
// array with promises of value
static any<R>(values: Promise.Thenable<R>[]): Promise<R>;
// array with values
static any<R>(values: R[]): Promise<R>;
/**
* Given an array, or a promise of an array, which contains promises (or a mix of promises and values) return a promise that is fulfilled or rejected as soon as a promise in the array is fulfilled or rejected with the respective rejection reason or fulfillment value.
*
* **Note** If you pass empty array or a sparse array with no values, or a promise/thenable for such, it will be forever pending.
*/
// promise of array with promises of value
static race<R>(values: Promise.Thenable<Promise.Thenable<R>[]>): Promise<R>;
// promise of array with values
static race<R>(values: Promise.Thenable<R[]>): Promise<R>;
// array with promises of value
static race<R>(values: Promise.Thenable<R>[]): Promise<R>;
// array with values
static race<R>(values: R[]): Promise<R>;
/**
* Initiate a competetive race between multiple promises or values (values will become immediately fulfilled promises). When `count` amount of promises have been fulfilled, the returned promise is fulfilled with an array that contains the fulfillment values of the winners in order of resolution.
*
* If too many promises are rejected so that the promise can never become fulfilled, it will be immediately rejected with an array of rejection reasons in the order they were thrown in.
*
* *The original array is not modified.*
*/
// promise of array with promises of value
static some<R>(values: Promise.Thenable<Promise.Thenable<R>[]>, count: number): Promise<R[]>;
// promise of array with values
static some<R>(values: Promise.Thenable<R[]>, count: number): Promise<R[]>;
// array with promises of value
static some<R>(values: Promise.Thenable<R>[], count: number): Promise<R[]>;
// array with values
static some<R>(values: R[], count: number): Promise<R[]>;
/**
* Like `Promise.all()` but instead of having to pass an array, the array is generated from the passed variadic arguments.
*/
// variadic array with promises of value
static join<R>(...values: Promise.Thenable<R>[]): Promise<R[]>;
// variadic array with values
static join<R>(...values: R[]): Promise<R[]>;
/**
* Map an array, or a promise of an array, which contains a promises (or a mix of promises and values) with the given `mapper` function with the signature `(item, index, arrayLength)` where `item` is the resolved value of a respective promise in the input array. If any promise in the input array is rejected the returned promise is rejected as well.
*
* If the `mapper` function returns promises or thenables, the returned promise will wait for all the mapped results to be resolved as well.
*
* *The original array is not modified.*
*/
// promise of array with promises of value
static map<R, U>(values: Promise.Thenable<Promise.Thenable<R>[]>, mapper: (item: R, index: number, arrayLength: number) => Promise.Thenable<U>): Promise<U[]>;
static map<R, U>(values: Promise.Thenable<Promise.Thenable<R>[]>, mapper: (item: R, index: number, arrayLength: number) => U): Promise<U[]>;
// promise of array with values
static map<R, U>(values: Promise.Thenable<R[]>, mapper: (item: R, index: number, arrayLength: number) => Promise.Thenable<U>): Promise<U[]>;
static map<R, U>(values: Promise.Thenable<R[]>, mapper: (item: R, index: number, arrayLength: number) => U): Promise<U[]>;
// array with promises of value
static map<R, U>(values: Promise.Thenable<R>[], mapper: (item: R, index: number, arrayLength: number) => Promise.Thenable<U>): Promise<U[]>;
static map<R, U>(values: Promise.Thenable<R>[], mapper: (item: R, index: number, arrayLength: number) => U): Promise<U[]>;
// array with values
static map<R, U>(values: R[], mapper: (item: R, index: number, arrayLength: number) => Promise.Thenable<U>): Promise<U[]>;
static map<R, U>(values: R[], mapper: (item: R, index: number, arrayLength: number) => U): Promise<U[]>;
/**
* Reduce an array, or a promise of an array, which contains a promises (or a mix of promises and values) with the given `reducer` function with the signature `(total, current, index, arrayLength)` where `item` is the resolved value of a respective promise in the input array. If any promise in the input array is rejected the returned promise is rejected as well.
*
* If the reducer function returns a promise or a thenable, the result for the promise is awaited for before continuing with next iteration.
*
* *The original array is not modified. If no `intialValue` is given and the array doesn't contain at least 2 items, the callback will not be called and `undefined` is returned. If `initialValue` is given and the array doesn't have at least 1 item, `initialValue` is returned.*
*/
// promise of array with promises of value
static reduce<R, U>(values: Promise.Thenable<Promise.Thenable<R>[]>, reducer: (total: U, current: R, index: number, arrayLength: number) => Promise.Thenable<U>, initialValue?: U): Promise<U>;
static reduce<R, U>(values: Promise.Thenable<Promise.Thenable<R>[]>, reducer: (total: U, current: R, index: number, arrayLength: number) => U, initialValue?: U): Promise<U>;
// promise of array with values
static reduce<R, U>(values: Promise.Thenable<R[]>, reducer: (total: U, current: R, index: number, arrayLength: number) => Promise.Thenable<U>, initialValue?: U): Promise<U>;
static reduce<R, U>(values: Promise.Thenable<R[]>, reducer: (total: U, current: R, index: number, arrayLength: number) => U, initialValue?: U): Promise<U>;
// array with promises of value
static reduce<R, U>(values: Promise.Thenable<R>[], reducer: (total: U, current: R, index: number, arrayLength: number) => Promise.Thenable<U>, initialValue?: U): Promise<U>;
static reduce<R, U>(values: Promise.Thenable<R>[], reducer: (total: U, current: R, index: number, arrayLength: number) => U, initialValue?: U): Promise<U>;
// array with values
static reduce<R, U>(values: R[], reducer: (total: U, current: R, index: number, arrayLength: number) => Promise.Thenable<U>, initialValue?: U): Promise<U>;
static reduce<R, U>(values: R[], reducer: (total: U, current: R, index: number, arrayLength: number) => U, initialValue?: U): Promise<U>;
/**
* Filter an array, or a promise of an array, which contains a promises (or a mix of promises and values) with the given `filterer` function with the signature `(item, index, arrayLength)` where `item` is the resolved value of a respective promise in the input array. If any promise in the input array is rejected the returned promise is rejected as well.
*
* The return values from the filtered functions are coerced to booleans, with the exception of promises and thenables which are awaited for their eventual result.
*
* *The original array is not modified.
*/
// promise of array with promises of value
static filter<R>(values: Promise.Thenable<Promise.Thenable<R>[]>, filterer: (item: R, index: number, arrayLength: number) => Promise.Thenable<boolean>): Promise<R[]>;
static filter<R>(values: Promise.Thenable<Promise.Thenable<R>[]>, filterer: (item: R, index: number, arrayLength: number) => boolean): Promise<R[]>;
// promise of array with values
static filter<R>(values: Promise.Thenable<R[]>, filterer: (item: R, index: number, arrayLength: number) => Promise.Thenable<boolean>): Promise<R[]>;
static filter<R>(values: Promise.Thenable<R[]>, filterer: (item: R, index: number, arrayLength: number) => boolean): Promise<R[]>;
// array with promises of value
static filter<R>(values: Promise.Thenable<R>[], filterer: (item: R, index: number, arrayLength: number) => Promise.Thenable<boolean>): Promise<R[]>;
static filter<R>(values: Promise.Thenable<R>[], filterer: (item: R, index: number, arrayLength: number) => boolean): Promise<R[]>;
// array with values
static filter<R>(values: R[], filterer: (item: R, index: number, arrayLength: number) => Promise.Thenable<boolean>): Promise<R[]>;
static filter<R>(values: R[], filterer: (item: R, index: number, arrayLength: number) => boolean): Promise<R[]>;
}
declare module Promise {
export interface RangeError extends Error {
}
export interface CancellationError extends Error {
}
export interface TimeoutError extends Error {
}
export interface TypeError extends Error {
}
export interface RejectionError extends Error {
}
export interface OperationalError extends Error {
}
// Ideally, we'd define e.g. "export class RangeError extends Error {}",
// but as Error is defined as an interface (not a class), TypeScript doesn't
// allow extending Error, only implementing it.
// However, if we want to catch() only a specific error type, we need to pass
// a constructor function to it. So, as a workaround, we define them here as such.
export function RangeError(): RangeError;
export function CancellationError(): CancellationError;
export function TimeoutError(): TimeoutError;
export function TypeError(): TypeError;
export function RejectionError(): RejectionError;
export function OperationalError(): OperationalError;
export interface Thenable<R> {
then<U>(onFulfilled: (value: R) => Thenable<U>, onRejected: (error: any) => Thenable<U>): Thenable<U>;
then<U>(onFulfilled: (value: R) => Thenable<U>, onRejected?: (error: any) => U): Thenable<U>;
then<U>(onFulfilled: (value: R) => U, onRejected: (error: any) => Thenable<U>): Thenable<U>;
then<U>(onFulfilled?: (value: R) => U, onRejected?: (error: any) => U): Thenable<U>;
}
export interface Resolver<R> {
/**
* Returns a reference to the controlled promise that can be passed to clients.
*/
promise: Promise<R>;
/**
* Resolve the underlying promise with `value` as the resolution value. If `value` is a thenable or a promise, the underlying promise will assume its state.
*/
resolve(value: R): void;
resolve(): void;
/**
* Reject the underlying promise with `reason` as the rejection reason.
*/
reject(reason: any): void;
/**
* Progress the underlying promise with `value` as the progression value.
*/
progress(value: any): void;
/**
* Gives you a callback representation of the `PromiseResolver`. Note that this is not a method but a property. The callback accepts error object in first argument and success values on the 2nd parameter and the rest, I.E. node js conventions.
*
* If the the callback is called with multiple success values, the resolver fullfills its promise with an array of the values.
*/
// TODO specify resolver callback
callback: (err: any, value: R, ...values: R[]) => void;
}
export interface Inspection<R> {
/**
* See if the underlying promise was fulfilled at the creation time of this inspection object.
*/
isFulfilled(): boolean;
/**
* See if the underlying promise was rejected at the creation time of this inspection object.
*/
isRejected(): boolean;
/**
* See if the underlying promise was defer at the creation time of this inspection object.
*/
isPending(): boolean;
/**
* Get the fulfillment value of the underlying promise. Throws if the promise wasn't fulfilled at the creation time of this inspection object.
*
* throws `TypeError`
*/
value(): R;
/**
* Get the rejection reason for the underlying promise. Throws if the promise wasn't rejected at the creation time of this inspection object.
*
* throws `TypeError`
*/
reason(): any;
}
/**
* Changes how bluebird schedules calls a-synchronously.
*
* @param scheduler Should be a function that asynchronously schedules
* the calling of the passed in function
*/
export function setScheduler(scheduler: (callback: (...args: any[]) => void) => void): void;
}
declare module 'bluebird' {
export = Promise;
}
@@ -1 +1 @@
/// <reference path="node/node.d.ts" />
/// <reference path="node/node.d.ts" />
@@ -164,7 +164,6 @@ public final class BuiltInWebServer extends HttpRequestHandler {
}
final String path = FileUtil.toCanonicalPath(decodedPath.substring(offset + 1), '/');
for (WebServerPathHandler pathHandler : WebServerPathHandler.EP_NAME.getExtensions()) {
try {
if (pathHandler.process(path, project, request, context, projectName, decodedPath, isCustomHost)) {
@@ -98,7 +98,6 @@ final class DefaultWebServerPathHandler extends WebServerPathHandler {
BuiltInWebServer.LOG.error(e);
}
}
return false;
}
}
@@ -1,6 +1,5 @@
package org.jetbrains.builtInWebServer;
import com.intellij.concurrency.JobScheduler;
import com.intellij.execution.ExecutionException;
import com.intellij.execution.filters.TextConsoleBuilder;
import com.intellij.execution.process.OSProcessHandler;
@@ -13,19 +12,17 @@ import com.intellij.openapi.actionSystem.DefaultActionGroup;
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.util.Consumer;
import com.intellij.util.net.NetUtils;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.concurrency.AsyncPromise;
import org.jetbrains.io.NettyUtil;
import org.jetbrains.concurrency.AsyncValueLoader;
import org.jetbrains.concurrency.Promise;
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);
@@ -50,19 +47,23 @@ public abstract class NetService implements Disposable {
}
}
@NotNull
@Override
protected void load(@NotNull final AsyncResult<OSProcessHandler> result) throws IOException {
protected Promise<OSProcessHandler> load(@NotNull final AsyncPromise<OSProcessHandler> promise) throws IOException {
final int port = NetUtils.findAvailableSocketPort();
final OSProcessHandler processHandler = doGetProcessHandler(port);
if (processHandler == null) {
result.setRejected();
return;
promise.setError(Promise.createError("rejected"));
return promise;
}
result.doWhenRejected(new Runnable() {
promise.rejected(new Consumer<Throwable>() {
@Override
public void run() {
public void consume(Throwable error) {
processHandler.destroyProcess();
if (!(error instanceof Promise.MessageError)) {
LOG.error(error);
}
}
});
@@ -70,35 +71,26 @@ public abstract class NetService implements Disposable {
processHandler.addProcessListener(processListener);
processHandler.startNotify();
if (result.isRejected()) {
return;
if (promise.getState() == Promise.State.REJECTED) {
return promise;
}
JobScheduler.getScheduler().schedule(new Runnable() {
ApplicationManager.getApplication().executeOnPooledThread(new Runnable() {
@Override
public void run() {
if (result.isRejected()) {
return;
}
ApplicationManager.getApplication().executeOnPooledThread(new Runnable() {
@Override
public void run() {
if (!result.isRejected()) {
AsyncPromise<OSProcessHandler> promise = new AsyncPromise<OSProcessHandler>();
try {
connectToProcess(promise, port, processHandler, processListener);
promise.notify(result);
}
catch (Throwable e) {
result.setRejected();
LOG.error(e);
}
if (promise.getState() != Promise.State.REJECTED) {
try {
connectToProcess(promise, port, processHandler, processListener);
}
catch (Throwable e) {
if (!promise.setError(e)) {
LOG.error(e);
}
}
});
}
}
}, NettyUtil.MIN_START_TIME, TimeUnit.MILLISECONDS);
});
return promise;
}
@Override
@@ -25,8 +25,9 @@ public abstract class SingleConnectionNetService extends NetService {
protected void connectToProcess(@NotNull AsyncPromise<OSProcessHandler> promise, int port, @NotNull OSProcessHandler processHandler, @NotNull Consumer<String> errorOutputConsumer) {
Bootstrap bootstrap = NettyUtil.oioClientBootstrap();
configureBootstrap(bootstrap, errorOutputConsumer);
processChannel = NettyUtil.connect(bootstrap, new InetSocketAddress(NetUtils.getLoopbackAddress(), port), promise);
if (processChannel != null) {
Channel channel = NettyUtil.connect(bootstrap, new InetSocketAddress(NetUtils.getLoopbackAddress(), port), promise);
if (channel != null) {
processChannel = channel;
promise.setResult(processHandler);
}
}
@@ -31,7 +31,7 @@ import java.net.URI;
* By default {@link WebServerPathToFileManager} will be used to map request to file.
* If file physically exists in the file system, you must use {@link WebServerRootsProvider}.
*
* Consider to extend {@link WebServerPathHandlerAdapter} instead of implement low-level {@link #process(String, com.intellij.openapi.project.Project, io.netty.handler.codec.http.FullHttpRequest, io.netty.channel.Channel, String, String, boolean)}
* Consider to extend {@link WebServerPathHandlerAdapter} instead of implement low-level {@link #process)}
*/
public abstract class WebServerPathHandler {
static final ExtensionPointName<WebServerPathHandler> EP_NAME = ExtensionPointName.create("org.jetbrains.webServerPathHandler");
@@ -46,9 +46,9 @@ public abstract class WebServerPathHandler {
protected static void redirectToDirectory(@NotNull HttpRequest request, @NotNull Channel channel, @NotNull String path) {
FullHttpResponse response = Responses.response(HttpResponseStatus.MOVED_PERMANENTLY);
URI url = VfsUtil.toUri("http://" + HttpHeaders.getHost(request) + '/' + path + '/');
URI url = VfsUtil.toUri("http://" + request.headers().get(HttpHeaderNames.HOST) + '/' + path + '/');
BuiltInWebServer.LOG.assertTrue(url != null);
response.headers().add(HttpHeaders.Names.LOCATION, url.toASCIIString());
response.headers().add(HttpHeaderNames.LOCATION, url.toASCIIString());
Responses.send(response, channel, request);
}
@@ -104,7 +104,7 @@ public class XmlRpcServerImpl implements XmlRpcServer {
}
Object response = invokeHandler(getHandler(xmlRpcServerRequest.getMethodName(), handlers == null ? handlerMapping : handlers), xmlRpcServerRequest);
result = Unpooled.copiedBuffer(new XmlRpcResponseProcessor().encodeResponse(response, CharsetToolkit.UTF8));
result = Unpooled.wrappedBuffer(new XmlRpcResponseProcessor().encodeResponse(response, CharsetToolkit.UTF8));
}
catch (Throwable e) {
context.channel().close();
@@ -74,7 +74,7 @@ public class BuiltInServer implements Disposable {
bootstrap.childHandler(new ChannelInitializer() {
@Override
protected void initChannel(Channel channel) throws Exception {
channel.pipeline().addLast(channelRegistrar, portUnificationServerHandler, ChannelExceptionHandler.getInstance());
channel.pipeline().addLast(channelRegistrar, portUnificationServerHandler);
}
});
}
@@ -85,7 +85,7 @@ public class BuiltInServer implements Disposable {
protected void initChannel(Channel channel) throws Exception {
channel.pipeline().addLast(channelRegistrar);
NettyUtil.addHttpServerCodec(channel.pipeline());
channel.pipeline().addLast(handler, ChannelExceptionHandler.getInstance());
channel.pipeline().addLast(handler);
}
});
}
@@ -175,7 +175,7 @@ public class BuiltInServer implements Disposable {
}
@Override
protected boolean process(ChannelHandlerContext context, FullHttpRequest request, QueryStringDecoder urlDecoder) throws IOException {
protected boolean process(@NotNull ChannelHandlerContext context, @NotNull FullHttpRequest request, @NotNull QueryStringDecoder urlDecoder) throws IOException {
if (handlers.isEmpty()) {
// not yet initialized, for example, P2PTransport could add handlers after we bound.
return false;
@@ -1,69 +0,0 @@
/*
* Copyright 2000-2014 JetBrains s.r.o.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.jetbrains.io;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
public abstract class Decoder extends ChannelInboundHandlerAdapter {
protected ByteBuf cumulation;
@Override
public final void channelRead(ChannelHandlerContext context, Object message) throws Exception {
if (message instanceof ByteBuf) {
messageReceived(context, (ByteBuf)message);
}
else {
context.fireChannelRead(message);
}
}
protected abstract void messageReceived(@NotNull ChannelHandlerContext context, @NotNull ByteBuf message) throws Exception;
@Nullable
protected final ByteBuf getBufferIfSufficient(@NotNull ByteBuf input, int requiredLength, @NotNull ChannelHandlerContext context) {
if (!input.isReadable()) {
return null;
}
if (cumulation == null) {
if (input.readableBytes() < requiredLength) {
cumulation = context.channel().config().getAllocator().buffer(requiredLength);
cumulation.writeBytes(input);
return null;
}
else {
return input;
}
}
else {
if ((cumulation.readableBytes() + input.readableBytes()) < requiredLength) {
cumulation.writeBytes(input);
return null;
}
else {
ByteBuf buffer = Unpooled.wrappedBuffer(cumulation, input);
input.skipBytes(input.readableBytes());
cumulation = null;
return buffer;
}
}
}
}
@@ -1,5 +1,5 @@
/*
* Copyright 2000-2013 JetBrains s.r.o.
* Copyright 2000-2015 JetBrains s.r.o.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -29,6 +29,7 @@ import io.netty.util.AttributeKey;
import org.apache.sanselan.ImageFormat;
import org.apache.sanselan.ImageWriteException;
import org.apache.sanselan.Sanselan;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.ide.HttpRequestHandler;
import javax.swing.*;
@@ -40,7 +41,7 @@ final class DelegatingHttpRequestHandler extends DelegatingHttpRequestHandlerBas
private static final AttributeKey<HttpRequestHandler> PREV_HANDLER = AttributeKey.valueOf("DelegatingHttpRequestHandler.handler");
@Override
protected boolean process(ChannelHandlerContext context, FullHttpRequest request, QueryStringDecoder urlDecoder) throws IOException, ImageWriteException {
protected boolean process(@NotNull ChannelHandlerContext context, @NotNull FullHttpRequest request, @NotNull QueryStringDecoder urlDecoder) throws IOException, ImageWriteException {
Attribute<HttpRequestHandler> prevHandlerAttribute = context.attr(PREV_HANDLER);
HttpRequestHandler connectedHandler = prevHandlerAttribute.get();
if (connectedHandler != null) {
@@ -80,9 +81,12 @@ final class DelegatingHttpRequestHandler extends DelegatingHttpRequestHandlerBas
}
@Override
public void exceptionCaught(ChannelHandlerContext context, Throwable cause) throws Exception {
super.exceptionCaught(context, cause);
context.attr(PREV_HANDLER).remove();
public void exceptionCaught(@NotNull ChannelHandlerContext context, @NotNull Throwable cause) {
try {
context.attr(PREV_HANDLER).remove();
}
finally {
super.exceptionCaught(context, cause);
}
}
}
@@ -1,5 +1,5 @@
/*
* Copyright 2000-2013 JetBrains s.r.o.
* Copyright 2000-2015 JetBrains s.r.o.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -19,12 +19,13 @@ import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.http.FullHttpRequest;
import io.netty.handler.codec.http.HttpResponseStatus;
import io.netty.handler.codec.http.QueryStringDecoder;
import org.jetbrains.annotations.NotNull;
abstract class DelegatingHttpRequestHandlerBase extends SimpleChannelInboundHandlerAdapter<FullHttpRequest> {
@Override
protected void messageReceived(ChannelHandlerContext context, FullHttpRequest message) throws Exception {
if (BuiltInServer.LOG.isDebugEnabled()) {
// BuiltInServer.LOG.debug("IN HTTP:\n" + message);
//BuiltInServer.LOG.debug("IN HTTP:\n" + message);
BuiltInServer.LOG.debug("IN HTTP: " + message.uri());
}
@@ -33,15 +34,10 @@ abstract class DelegatingHttpRequestHandlerBase extends SimpleChannelInboundHand
}
}
protected abstract boolean process(ChannelHandlerContext context, FullHttpRequest request, QueryStringDecoder urlDecoder) throws Exception;
protected abstract boolean process(@NotNull ChannelHandlerContext context, @NotNull FullHttpRequest request, @NotNull QueryStringDecoder urlDecoder) throws Exception;
@Override
public void exceptionCaught(ChannelHandlerContext context, Throwable cause) throws Exception {
try {
NettyUtil.log(cause, BuiltInServer.LOG);
}
finally {
context.channel().close();
}
public void exceptionCaught(@NotNull ChannelHandlerContext context, @NotNull Throwable cause) {
NettyUtil.logAndClose(cause, BuiltInServer.LOG, context.channel());
}
}
@@ -45,10 +45,10 @@ public class FileResponses {
}
private static boolean checkCache(@NotNull HttpRequest request, @NotNull Channel channel, long lastModified) {
String ifModifiedSince = request.headers().get(HttpHeaders.Names.IF_MODIFIED_SINCE);
String ifModifiedSince = request.headers().get(HttpHeaderNames.IF_MODIFIED_SINCE);
if (!StringUtil.isEmpty(ifModifiedSince)) {
try {
if (HttpHeaders.getDateHeader(request, HttpHeaders.Names.IF_MODIFIED_SINCE).getTime() >= lastModified) {
if (HttpHeaders.getDateHeader(request, HttpHeaderNames.IF_MODIFIED_SINCE).getTime() >= lastModified) {
send(response(HttpResponseStatus.NOT_MODIFIED), channel, request);
return true;
}
@@ -70,8 +70,8 @@ public class FileResponses {
HttpResponse response = new DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK);
response.headers().add(CONTENT_TYPE, getContentType(path));
addCommonHeaders(response);
response.headers().set(HttpHeaders.Names.CACHE_CONTROL, "private, must-revalidate");
HttpHeaders.setDateHeader(response, HttpHeaders.Names.LAST_MODIFIED, new Date(lastModified));
response.headers().set(HttpHeaderNames.CACHE_CONTROL, "private, must-revalidate");
HttpHeaders.setDateHeader(response, HttpHeaderNames.LAST_MODIFIED, new Date(lastModified));
return response;
}
@@ -74,19 +74,16 @@ class PortUnificationServerHandler extends Decoder {
this(new DelegatingHttpRequestHandler(), true, true);
}
private PortUnificationServerHandler(DelegatingHttpRequestHandler delegatingHttpRequestHandler, boolean detectSsl, boolean detectGzip) {
private PortUnificationServerHandler(@NotNull DelegatingHttpRequestHandler delegatingHttpRequestHandler, boolean detectSsl, boolean detectGzip) {
this.delegatingHttpRequestHandler = delegatingHttpRequestHandler;
this.detectSsl = detectSsl;
this.detectGzip = detectGzip;
}
@Override
protected void messageReceived(@NotNull ChannelHandlerContext context, @NotNull ByteBuf message) throws Exception {
ByteBuf buffer = getBufferIfSufficient(message, 5, context);
if (buffer == null) {
message.release();
}
else {
protected void messageReceived(@NotNull ChannelHandlerContext context, @NotNull ByteBuf input) throws Exception {
ByteBuf buffer = getBufferIfSufficient(input, 5, context);
if (buffer != null) {
decode(context, buffer);
}
}
@@ -108,11 +105,7 @@ class PortUnificationServerHandler extends Decoder {
}
else if (isHttp(magic1, magic2)) {
NettyUtil.addHttpServerCodec(pipeline);
pipeline.addLast(delegatingHttpRequestHandler);
// added earlier if HTTPS
if (pipeline.get(ChunkedWriteHandler.class) == null) {
pipeline.addLast("chunkedWriteHandler", new ChunkedWriteHandler());
}
pipeline.addLast("delegatingHttpHandler", delegatingHttpRequestHandler);
if (BuiltInServer.LOG.isDebugEnabled()) {
pipeline.addLast(new ChannelOutboundHandlerAdapter() {
@Override
@@ -139,17 +132,17 @@ class PortUnificationServerHandler extends Decoder {
// must be after new channels handlers addition (netty bug?)
pipeline.remove(this);
ensureThatExceptionHandlerIsLast(pipeline);
// Buffer will be automatically released after messageReceived, but we pass it to next handler, and next handler will also release, so, we must retain.
// We can introduce Decoder.isAutoRelease, but in this case, if error will be thrown while we are executing, buffer will not be released.
// So, it is robust solution just always release (Decoder does) and just retain (we - client) if autorelease behavior is not suitable.
buffer.retain();
// we must fire channel read - new added handler must read buffer
context.fireChannelRead(buffer);
}
private static void ensureThatExceptionHandlerIsLast(ChannelPipeline pipeline) {
ChannelHandler exceptionHandler = ChannelExceptionHandler.getInstance();
if (pipeline.last() != exceptionHandler) {
pipeline.remove(exceptionHandler);
pipeline.addLast(exceptionHandler);
}
@Override
public void exceptionCaught(ChannelHandlerContext context, Throwable cause) throws Exception {
NettyUtil.logAndClose(cause, BuiltInServer.LOG, context.channel());
}
private static boolean isHttp(int magic1, int magic2) {
@@ -169,24 +162,29 @@ class PortUnificationServerHandler extends Decoder {
private static final int UUID_LENGTH = 16;
@Override
protected void messageReceived(@NotNull ChannelHandlerContext context, @NotNull ByteBuf message) throws Exception {
ByteBuf buffer = getBufferIfSufficient(message, UUID_LENGTH, context);
protected void messageReceived(@NotNull ChannelHandlerContext context, @NotNull ByteBuf input) throws Exception {
ByteBuf buffer = getBufferIfSufficient(input, UUID_LENGTH, context);
if (buffer == null) {
message.release();
return;
}
else {
UUID uuid = new UUID(buffer.readLong(), buffer.readLong());
for (BinaryRequestHandler customHandler : BinaryRequestHandler.EP_NAME.getExtensions()) {
if (uuid.equals(customHandler.getId())) {
ChannelPipeline pipeline = context.pipeline();
pipeline.addLast(customHandler.getInboundHandler());
pipeline.remove(this);
ensureThatExceptionHandlerIsLast(pipeline);
context.fireChannelRead(buffer);
break;
}
UUID uuid = new UUID(buffer.readLong(), buffer.readLong());
for (BinaryRequestHandler customHandler : BinaryRequestHandler.EP_NAME.getExtensions()) {
if (uuid.equals(customHandler.getId())) {
ChannelPipeline pipeline = context.pipeline();
pipeline.addLast(customHandler.getInboundHandler(context));
pipeline.addLast(ChannelExceptionHandler.getInstance());
pipeline.remove(this);
context.fireChannelRead(buffer);
break;
}
}
}
@Override
public void exceptionCaught(ChannelHandlerContext context, Throwable cause) {
NettyUtil.logAndClose(cause, BuiltInServer.LOG, context.channel());
}
}
}
@@ -20,6 +20,8 @@ import com.intellij.openapi.application.ApplicationManager;
import com.intellij.openapi.application.ex.ApplicationInfoEx;
import com.intellij.openapi.util.text.StringUtil;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.ByteBufAllocator;
import io.netty.buffer.ByteBufUtil;
import io.netty.buffer.Unpooled;
import io.netty.channel.Channel;
import io.netty.channel.ChannelFuture;
@@ -29,6 +31,7 @@ import io.netty.util.CharsetUtil;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.nio.CharBuffer;
import java.nio.charset.Charset;
import java.util.Calendar;
@@ -45,20 +48,20 @@ public final class Responses {
? new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK)
: new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK, content);
if (contentType != null) {
response.headers().set(HttpHeaders.Names.CONTENT_TYPE, contentType);
response.headers().set(HttpHeaderNames.CONTENT_TYPE, contentType);
}
return response;
}
public static void setDate(@NotNull HttpResponse response) {
if (!response.headers().contains(HttpHeaders.Names.DATE)) {
HttpHeaders.setDateHeader(response, HttpHeaders.Names.DATE, Calendar.getInstance().getTime());
if (!response.headers().contains(HttpHeaderNames.DATE)) {
HttpHeaders.setDateHeader(response, HttpHeaderNames.DATE, Calendar.getInstance().getTime());
}
}
public static void addNoCache(@NotNull HttpResponse response) {
response.headers().add(HttpHeaders.Names.CACHE_CONTROL, "no-cache, no-store, must-revalidate, max-age=0");
response.headers().add(HttpHeaders.Names.PRAGMA, "no-cache");
response.headers().add(HttpHeaderNames.CACHE_CONTROL, "no-cache, no-store, must-revalidate, max-age=0");
response.headers().add(HttpHeaderNames.PRAGMA, "no-cache");
}
@Nullable
@@ -74,13 +77,13 @@ public final class Responses {
public static void addServer(@NotNull HttpResponse response) {
if (getServerHeaderValue() != null) {
response.headers().add(HttpHeaders.Names.SERVER, getServerHeaderValue());
response.headers().add(HttpHeaderNames.SERVER, getServerHeaderValue());
}
}
public static void send(@NotNull HttpResponse response, Channel channel, @Nullable HttpRequest request) {
if (response.status() != HttpResponseStatus.NOT_MODIFIED && !HttpHeaders.isContentLengthSet(response)) {
HttpHeaders.setContentLength(response,
if (response.status() != HttpResponseStatus.NOT_MODIFIED && !HttpHeaderUtil.isContentLengthSet(response)) {
HttpHeaderUtil.setContentLength(response,
response instanceof FullHttpResponse ? ((FullHttpResponse)response).content().readableBytes() : 0);
}
@@ -89,8 +92,8 @@ public final class Responses {
}
public static boolean addKeepAliveIfNeed(HttpResponse response, HttpRequest request) {
if (HttpHeaders.isKeepAlive(request)) {
HttpHeaders.setKeepAlive(response, true);
if (HttpHeaderUtil.isKeepAlive(request)) {
HttpHeaderUtil.setKeepAlive(response, true);
return true;
}
return false;
@@ -149,8 +152,8 @@ public final class Responses {
}
builder.append("<hr/><p style=\"text-align: center\">").append(StringUtil.notNullize(getServerHeaderValue(), "")).append("</p>");
DefaultFullHttpResponse response = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, responseStatus, Unpooled.copiedBuffer(builder, CharsetUtil.UTF_8));
response.headers().set(HttpHeaders.Names.CONTENT_TYPE, "text/html");
DefaultFullHttpResponse response = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, responseStatus, ByteBufUtil.encodeString(ByteBufAllocator.DEFAULT, CharBuffer.wrap(builder), CharsetUtil.UTF_8));
response.headers().set(HttpHeaderNames.CONTENT_TYPE, "text/html");
return response;
}
}
@@ -1,116 +0,0 @@
package org.jetbrains.io.fastCgi;
import com.intellij.openapi.util.text.StringUtil;
import com.intellij.openapi.util.text.StringUtilRt;
import com.intellij.util.containers.ConcurrentIntObjectMap;
import io.netty.buffer.ByteBuf;
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;
import static org.jetbrains.io.fastCgi.FastCgiService.LOG;
@ChannelHandler.Sharable
public class FastCgiChannelHandler extends SimpleChannelInboundHandlerAdapter<FastCgiResponse> {
private final ConcurrentIntObjectMap<Channel> requestToChannel;
public FastCgiChannelHandler(@NotNull ConcurrentIntObjectMap<Channel> channel) {
requestToChannel = channel;
}
@Override
protected void messageReceived(ChannelHandlerContext context, FastCgiResponse response) throws Exception {
ByteBuf buffer = response.getData();
Channel channel = requestToChannel.remove(response.getId());
if (channel == null || !channel.isActive()) {
if (buffer != null) {
buffer.release();
}
return;
}
if (buffer == null) {
Responses.sendStatus(HttpResponseStatus.BAD_GATEWAY, channel);
return;
}
HttpResponse httpResponse = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK, buffer);
try {
parseHeaders(httpResponse, buffer);
Responses.addServer(httpResponse);
if (!HttpHeaders.isContentLengthSet(httpResponse)) {
HttpHeaders.setContentLength(httpResponse, buffer.readableBytes());
}
}
catch (Throwable e) {
buffer.release();
Responses.sendStatus(HttpResponseStatus.INTERNAL_SERVER_ERROR, channel);
LOG.error(e);
}
channel.writeAndFlush(httpResponse);
}
private static void parseHeaders(HttpResponse response, ByteBuf buffer) {
StringBuilder builder = new StringBuilder();
while (buffer.isReadable()) {
builder.setLength(0);
String key = null;
boolean valueExpected = true;
while (true) {
int b = buffer.readByte();
if (b < 0 || b == '\n') {
break;
}
if (b != '\r') {
if (valueExpected && b == ':') {
valueExpected = false;
key = builder.toString();
builder.setLength(0);
skipWhitespace(buffer);
}
else {
builder.append((char)b);
}
}
}
if (builder.length() == 0) {
// end of headers
return;
}
// skip standard headers
if (StringUtil.isEmpty(key) || StringUtilRt.startsWithIgnoreCase(key, "http") || StringUtilRt.startsWithIgnoreCase(key, "X-Accel-")) {
continue;
}
String value = builder.toString();
if (key.equalsIgnoreCase("status")) {
int index = value.indexOf(' ');
if (index == -1) {
LOG.warn("Cannot parse status: " + value);
response.setStatus(HttpResponseStatus.OK);
}
else {
response.setStatus(HttpResponseStatus.valueOf(Integer.parseInt(value.substring(0, index))));
}
}
else if (!(key.startsWith("http") || key.startsWith("HTTP"))) {
response.headers().add(key, value);
}
}
}
private static void skipWhitespace(ByteBuf buffer) {
while (buffer.isReadable() && buffer.getByte(buffer.readerIndex()) == ' ') {
buffer.skipBytes(1);
}
}
}
@@ -2,9 +2,9 @@ package org.jetbrains.io.fastCgi;
import com.intellij.util.Consumer;
import gnu.trove.TIntObjectHashMap;
import gnu.trove.TIntObjectProcedure;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.CompositeByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandlerContext;
import io.netty.util.CharsetUtil;
import org.jetbrains.annotations.NotNull;
@@ -12,7 +12,7 @@ import org.jetbrains.io.Decoder;
import static org.jetbrains.io.fastCgi.FastCgiService.LOG;
public class FastCgiDecoder extends Decoder {
public class FastCgiDecoder extends Decoder implements Decoder.FullMessageConsumer<Void> {
private enum State {
HEADER, CONTENT
}
@@ -37,9 +37,11 @@ public class FastCgiDecoder extends Decoder {
private final TIntObjectHashMap<ByteBuf> dataBuffers = new TIntObjectHashMap<ByteBuf>();
private final Consumer<String> errorOutputConsumer;
private final FastCgiService responseHandler;
public FastCgiDecoder(Consumer<String> errorOutputConsumer) {
public FastCgiDecoder(@NotNull Consumer<String> errorOutputConsumer, @NotNull FastCgiService responseHandler) {
this.errorOutputConsumer = errorOutputConsumer;
this.responseHandler = responseHandler;
}
@Override
@@ -48,21 +50,19 @@ public class FastCgiDecoder extends Decoder {
switch (state) {
case HEADER: {
if (paddingLength > 0) {
if (input.readableBytes() >= paddingLength) {
if (input.readableBytes() > paddingLength) {
input.skipBytes(paddingLength);
paddingLength = 0;
}
else {
paddingLength -= input.readableBytes();
input.skipBytes(input.readableBytes());
input.release();
return;
}
}
ByteBuf buffer = getBufferIfSufficient(input, FastCgiConstants.HEADER_LENGTH, context);
if (buffer == null) {
input.release();
return;
}
@@ -72,16 +72,7 @@ public class FastCgiDecoder extends Decoder {
case CONTENT: {
if (contentLength > 0) {
ByteBuf buffer = getBufferIfSufficient(input, contentLength, context);
if (buffer == null) {
input.release();
return;
}
FastCgiResponse response = readContent(buffer);
if (response != null) {
context.fireChannelRead(response);
}
readContent(input, context, contentLength, this);
}
state = State.HEADER;
}
@@ -89,7 +80,31 @@ public class FastCgiDecoder extends Decoder {
}
}
private void decodeHeader(ByteBuf buffer) {
@Override
public void channelInactive(ChannelHandlerContext context) throws Exception {
try {
if (!dataBuffers.isEmpty()) {
dataBuffers.forEachEntry(new TIntObjectProcedure<ByteBuf>() {
@Override
public boolean execute(int a, ByteBuf buffer) {
try {
buffer.release();
}
catch (Throwable e) {
LOG.error(e);
}
return true;
}
});
dataBuffers.clear();
}
}
finally {
super.channelInactive(context);
}
}
private void decodeHeader(@NotNull ByteBuf buffer) {
buffer.skipBytes(1);
type = buffer.readUnsignedByte();
id = buffer.readUnsignedShort();
@@ -98,25 +113,25 @@ public class FastCgiDecoder extends Decoder {
buffer.skipBytes(1);
}
private FastCgiResponse readContent(ByteBuf buffer) {
@Override
public Void contentReceived(@NotNull ByteBuf buffer, @NotNull ChannelHandlerContext context, boolean isCumulateBuffer) {
switch (type) {
case RecordType.END_REQUEST:
int appStatus = buffer.readInt();
int protocolStatus = buffer.readUnsignedByte();
buffer.skipBytes(3);
if (appStatus != 0 || protocolStatus != ProtocolStatus.REQUEST_COMPLETE.ordinal()) {
LOG.warn("Protocol status " + protocolStatus);
dataBuffers.remove(id);
return new FastCgiResponse(id, null);
responseHandler.responseReceived(id, null);
}
else if (protocolStatus == ProtocolStatus.REQUEST_COMPLETE.ordinal()) {
return new FastCgiResponse(id, dataBuffers.remove(id));
responseHandler.responseReceived(id, dataBuffers.remove(id));
}
break;
case RecordType.STDOUT:
ByteBuf data = dataBuffers.get(id);
ByteBuf sliced = buffer.slice(buffer.readerIndex(), contentLength);
ByteBuf sliced = isCumulateBuffer ? buffer : buffer.slice(buffer.readerIndex(), contentLength);
if (data == null) {
dataBuffers.put(id, sliced);
}
@@ -125,10 +140,17 @@ public class FastCgiDecoder extends Decoder {
data.writerIndex(data.writerIndex() + sliced.readableBytes());
}
else {
dataBuffers.put(id, Unpooled.wrappedBuffer(data, sliced));
if (sliced instanceof CompositeByteBuf) {
data = ((CompositeByteBuf)sliced).addComponent(0, data);
data.writerIndex(data.writerIndex() + data.readableBytes());
}
else {
data = context.alloc().compositeBuffer(DEFAULT_MAX_COMPOSITE_BUFFER_COMPONENTS).addComponents(data, sliced);
data.writerIndex(data.writerIndex() + data.readableBytes() + sliced.readableBytes());
}
dataBuffers.put(id, data);
}
sliced.retain();
buffer.skipBytes(contentLength);
break;
case RecordType.STDERR:
@@ -138,7 +160,6 @@ public class FastCgiDecoder extends Decoder {
catch (Throwable e) {
LOG.error(e);
}
buffer.skipBytes(contentLength);
break;
default:
@@ -4,11 +4,10 @@ import com.intellij.openapi.project.Project;
import com.intellij.openapi.vfs.VirtualFile;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.ByteBufAllocator;
import io.netty.buffer.Unpooled;
import io.netty.buffer.ByteBufUtil;
import io.netty.channel.Channel;
import io.netty.handler.codec.http.FullHttpRequest;
import io.netty.handler.codec.http.HttpHeaders;
import io.netty.util.CharsetUtil;
import io.netty.handler.codec.http.HttpHeaderNames;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.builtInWebServer.PathInfo;
@@ -27,16 +26,17 @@ public class FastCgiRequest {
private static final int STDIN = 5;
private static final int VERSION = 1;
private final ByteBuf buffer;
private ByteBuf buffer;
final int requestId;
public FastCgiRequest(int requestId, @NotNull ByteBufAllocator allocator) {
this.requestId = requestId;
buffer = allocator.buffer();
buffer = allocator.ioBuffer(4096);
writeHeader(buffer, BEGIN_REQUEST, FastCgiConstants.HEADER_LENGTH);
buffer.writeShort(RESPONDER);
buffer.writeByte(FCGI_KEEP_CONNECTION);
// reserved[5]
buffer.writeZero(5);
}
@@ -77,11 +77,11 @@ public class FastCgiRequest {
buffer.writeByte(valLength);
}
buffer.writeBytes(key.getBytes(CharsetUtil.US_ASCII));
buffer.writeBytes(Unpooled.copiedBuffer(value, CharsetUtil.UTF_8));
ByteBufUtil.writeAscii(buffer, key);
ByteBufUtil.writeUtf8(buffer, value);
}
public void writeHeaders(FullHttpRequest request, Channel clientChannel) {
public void writeHeaders(@NotNull FullHttpRequest request, @NotNull Channel clientChannel) {
addHeader("REQUEST_URI", request.uri());
addHeader("REQUEST_METHOD", request.method().name());
@@ -98,7 +98,7 @@ public class FastCgiRequest {
addHeader("GATEWAY_INTERFACE", "CGI/1.1");
addHeader("SERVER_PROTOCOL", request.protocolVersion().text());
addHeader("CONTENT_TYPE", request.headers().get(HttpHeaders.Names.CONTENT_TYPE));
addHeader("CONTENT_TYPE", request.headers().get(HttpHeaderNames.CONTENT_TYPE));
// PHP only, required if PHP was built with --enable-force-cgi-redirect
addHeader("REDIRECT_STATUS", "200");
@@ -117,33 +117,48 @@ public class FastCgiRequest {
}
}
final void writeToServerChannel(ByteBuf content, Channel fastCgiChannel) {
writeHeader(buffer, PARAMS, 0);
fastCgiChannel.write(buffer);
if (content.isReadable()) {
ByteBuf headerBuffer = fastCgiChannel.alloc().buffer(FastCgiConstants.HEADER_LENGTH, FastCgiConstants.HEADER_LENGTH);
writeHeader(headerBuffer, STDIN, content.readableBytes());
fastCgiChannel.write(headerBuffer);
fastCgiChannel.write(content);
headerBuffer = fastCgiChannel.alloc().buffer(FastCgiConstants.HEADER_LENGTH, FastCgiConstants.HEADER_LENGTH);
writeHeader(headerBuffer, STDIN, 0);
fastCgiChannel.write(headerBuffer);
final void writeToServerChannel(@Nullable ByteBuf content, @NotNull Channel fastCgiChannel) {
if (fastCgiChannel.pipeline().first() == null) {
throw new IllegalStateException("No handler in the pipeline");
}
else {
content.release();
boolean releaseContent = content != null;
try {
writeHeader(buffer, PARAMS, 0);
if (content != null) {
writeHeader(buffer, STDIN, content.readableBytes());
}
fastCgiChannel.write(buffer);
buffer = null;
if (content != null) {
fastCgiChannel.write(content);
// channel.write releases
releaseContent = false;
ByteBuf headerBuffer = fastCgiChannel.alloc().ioBuffer(FastCgiConstants.HEADER_LENGTH, FastCgiConstants.HEADER_LENGTH);
writeHeader(headerBuffer, STDIN, 0);
fastCgiChannel.write(headerBuffer);
}
}
finally {
if (releaseContent) {
assert content != null;
content.release();
}
}
fastCgiChannel.flush();
}
private void writeHeader(ByteBuf buffer, int type, int length) {
private void writeHeader(@NotNull ByteBuf buffer, int type, int length) {
buffer.writeByte(VERSION);
buffer.writeByte(type);
buffer.writeShort(requestId);
buffer.writeShort(length);
// paddingLength, reserved
buffer.writeZero(2);
}
}
@@ -1,21 +0,0 @@
package org.jetbrains.io.fastCgi;
import io.netty.buffer.ByteBuf;
public class FastCgiResponse {
private final int id;
private final ByteBuf data;
public FastCgiResponse(int id, ByteBuf data) {
this.id = id;
this.data = data;
}
public ByteBuf getData() {
return data;
}
public int getId() {
return id;
}
}
@@ -15,8 +15,11 @@
*/
package org.jetbrains.io.fastCgi;
import com.intellij.execution.process.OSProcessHandler;
import com.intellij.openapi.diagnostic.Logger;
import com.intellij.openapi.project.Project;
import com.intellij.openapi.util.text.StringUtil;
import com.intellij.openapi.util.text.StringUtilRt;
import com.intellij.util.Consumer;
import com.intellij.util.containers.ConcurrentIntObjectMap;
import com.intellij.util.containers.ContainerUtil;
@@ -24,10 +27,12 @@ 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 io.netty.handler.codec.http.*;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.builtInWebServer.SingleConnectionNetService;
import org.jetbrains.io.ChannelExceptionHandler;
import org.jetbrains.io.MessageDecoder;
import org.jetbrains.io.NettyUtil;
import org.jetbrains.io.Responses;
@@ -36,7 +41,7 @@ import java.util.concurrent.atomic.AtomicInteger;
// todo send FCGI_ABORT_REQUEST if client channel disconnected
public abstract class FastCgiService extends SingleConnectionNetService {
static final Logger LOG = Logger.getInstance(FastCgiService.class);
protected static final Logger LOG = Logger.getInstance(FastCgiService.class);
private final AtomicInteger requestIdCounter = new AtomicInteger();
protected final ConcurrentIntObjectMap<Channel> requests = ContainerUtil.createConcurrentIntObjectMap();
@@ -56,56 +61,87 @@ public abstract class FastCgiService extends SingleConnectionNetService {
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);
}
sendBadGateway(channel);
}
}
}
}
private static void sendBadGateway(@NotNull Channel channel) {
try {
if (channel.isActive()) {
Responses.sendStatus(HttpResponseStatus.BAD_GATEWAY, channel);
}
}
catch (Throwable e) {
NettyUtil.log(e, LOG);
}
}
@Override
protected void configureBootstrap(@NotNull Bootstrap bootstrap, @NotNull final Consumer<String> errorOutputConsumer) {
final FastCgiChannelHandler fastCgiChannelHandler = new FastCgiChannelHandler(requests);
bootstrap.handler(new ChannelInitializer() {
@Override
protected void initChannel(Channel channel) throws Exception {
channel.pipeline().addLast(new FastCgiDecoder(errorOutputConsumer), fastCgiChannelHandler, ChannelExceptionHandler.getInstance());
channel.pipeline().addLast("fastCgiDecoder", new FastCgiDecoder(errorOutputConsumer, FastCgiService.this));
channel.pipeline().addLast("exceptionHandler", ChannelExceptionHandler.getInstance());
}
});
}
public void send(final FastCgiRequest fastCgiRequest, final ByteBuf content) {
content.retain();
if (processHandler.has()) {
fastCgiRequest.writeToServerChannel(content, processChannel);
public void send(@NotNull final FastCgiRequest fastCgiRequest, @NotNull ByteBuf content) {
final ByteBuf notEmptyContent;
if (content.isReadable()) {
content.retain();
notEmptyContent = content;
notEmptyContent.touch();
}
else {
processHandler.get().doWhenDone(new Runnable() {
@Override
public void run() {
fastCgiRequest.writeToServerChannel(content, processChannel);
}
}).doWhenRejected(new Runnable() {
@Override
public void run() {
content.release();
Channel channel = requests.get(fastCgiRequest.requestId);
if (channel != null && channel.isActive()) {
Responses.sendStatus(HttpResponseStatus.BAD_GATEWAY, channel);
}
}
});
notEmptyContent = null;
}
try {
if (processHandler.has()) {
fastCgiRequest.writeToServerChannel(notEmptyContent, processChannel);
}
else {
processHandler.get()
.done(new Consumer<OSProcessHandler>() {
@Override
public void consume(OSProcessHandler osProcessHandler) {
fastCgiRequest.writeToServerChannel(notEmptyContent, processChannel);
}
})
.rejected(new Consumer<Throwable>() {
@Override
public void consume(Throwable error) {
LOG.error(error);
handleError(fastCgiRequest, notEmptyContent);
}
});
}
}
catch (Throwable e) {
LOG.error(e);
handleError(fastCgiRequest, notEmptyContent);
}
}
public int allocateRequestId(Channel channel) {
private void handleError(@NotNull FastCgiRequest fastCgiRequest, @Nullable ByteBuf content) {
try {
if (content != null && content.refCnt() != 0) {
content.release();
}
}
finally {
Channel channel = requests.remove(fastCgiRequest.requestId);
if (channel != null) {
sendBadGateway(channel);
}
}
}
public int allocateRequestId(@NotNull Channel channel) {
int requestId = requestIdCounter.getAndIncrement();
if (requestId >= Short.MAX_VALUE) {
requestIdCounter.set(0);
@@ -114,4 +150,93 @@ public abstract class FastCgiService extends SingleConnectionNetService {
requests.put(requestId, channel);
return requestId;
}
void responseReceived(int id, @Nullable ByteBuf buffer) {
Channel channel = requests.remove(id);
if (channel == null || !channel.isActive()) {
if (buffer != null) {
buffer.release();
}
return;
}
if (buffer == null) {
Responses.sendStatus(HttpResponseStatus.BAD_GATEWAY, channel);
return;
}
HttpResponse httpResponse = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK, buffer);
try {
parseHeaders(httpResponse, buffer);
Responses.addServer(httpResponse);
if (!HttpHeaderUtil.isContentLengthSet(httpResponse)) {
HttpHeaderUtil.setContentLength(httpResponse, buffer.readableBytes());
}
}
catch (Throwable e) {
buffer.release();
try {
LOG.error(e);
}
finally {
Responses.sendStatus(HttpResponseStatus.INTERNAL_SERVER_ERROR, channel);
}
return;
}
channel.writeAndFlush(httpResponse);
}
private static void parseHeaders(@NotNull HttpResponse response, @NotNull ByteBuf buffer) {
StringBuilder builder = new StringBuilder();
while (buffer.isReadable()) {
builder.setLength(0);
String key = null;
boolean valueExpected = true;
while (true) {
int b = buffer.readByte();
if (b < 0 || b == '\n') {
break;
}
if (b != '\r') {
if (valueExpected && b == ':') {
valueExpected = false;
key = builder.toString();
builder.setLength(0);
MessageDecoder.skipWhitespace(buffer);
}
else {
builder.append((char)b);
}
}
}
if (builder.length() == 0) {
// end of headers
return;
}
// skip standard headers
if (StringUtil.isEmpty(key) || StringUtilRt.startsWithIgnoreCase(key, "http") || StringUtilRt.startsWithIgnoreCase(key, "X-Accel-")) {
continue;
}
String value = builder.toString();
if (key.equalsIgnoreCase("status")) {
int index = value.indexOf(' ');
if (index == -1) {
LOG.warn("Cannot parse status: " + value);
response.setStatus(HttpResponseStatus.OK);
}
else {
response.setStatus(HttpResponseStatus.valueOf(Integer.parseInt(value.substring(0, index))));
}
}
else if (!(key.startsWith("http") || key.startsWith("HTTP"))) {
response.headers().add(key, value);
}
}
}
}
@@ -4,6 +4,7 @@ import com.intellij.openapi.util.UserDataHolderBase;
import com.intellij.util.containers.ConcurrentIntObjectMap;
import com.intellij.util.containers.ContainerUtil;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.ByteBufAllocator;
import io.netty.channel.Channel;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelFutureListener;
@@ -27,10 +28,15 @@ public abstract class Client extends UserDataHolderBase {
}
@NotNull
public EventLoop getEventLoop() {
public final EventLoop getEventLoop() {
return channel.eventLoop();
}
@NotNull
public final ByteBufAllocator getByteBufAllocator() {
return channel.alloc();
}
protected abstract ChannelFuture send(@NotNull ByteBuf message);
public abstract void sendHeartbeat();
@@ -49,7 +55,7 @@ public abstract class Client extends UserDataHolderBase {
return promise;
}
void rejectAsyncResults(@NotNull ExceptionHandler exceptionHandler) {
final void rejectAsyncResults(@NotNull ExceptionHandler exceptionHandler) {
if (!messageCallbackMap.isEmpty()) {
Enumeration<AsyncPromise<Object>> elements = messageCallbackMap.elements();
while (elements.hasMoreElements()) {
@@ -73,10 +79,10 @@ public abstract class Client extends UserDataHolderBase {
}
@Override
public void setError(@NotNull Throwable error) {
super.setError(error);
public boolean setError(@NotNull Throwable error) {
boolean result = super.setError(error);
messageCallbackMap.remove(messageId);
return result;
}
@SuppressWarnings("ThrowableResultOfMethodCallIgnored")
@@ -1,13 +1,14 @@
package org.jetbrains.io.jsonRpc;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.util.EventListener;
import java.util.List;
import java.util.Map;
public interface ClientListener extends EventListener {
void connected(@NotNull Client client, Map<String, List<String>> parameters);
void connected(@NotNull Client client, @Nullable Map<String, List<String>> parameters);
void disconnected(@NotNull Client client);
}
@@ -1,13 +1,14 @@
package org.jetbrains.io.jsonRpc;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.util.List;
import java.util.Map;
public abstract class ClientListenerAdapter implements ClientListener {
@Override
public void connected(@NotNull Client client, Map<String, List<String>> parameters) {
public void connected(@NotNull Client client, @Nullable Map<String, List<String>> parameters) {
}
@Override
@@ -1,8 +1,10 @@
package org.jetbrains.io.jsonRpc;
import org.jetbrains.annotations.NotNull;
public interface ExceptionHandler {
/**
* @param e Exception while encode message (on send)
*/
void exceptionCaught(Throwable e);
void exceptionCaught(@NotNull Throwable e);
}
@@ -1,8 +1,10 @@
package org.jetbrains.io.jsonRpc;
import org.jetbrains.annotations.NotNull;
public class ExceptionHandlerImpl implements ExceptionHandler {
@Override
public void exceptionCaught(Throwable e) {
public void exceptionCaught(@NotNull Throwable e) {
//noinspection CallToPrintStackTrace
e.printStackTrace();
}
@@ -0,0 +1,56 @@
package org.jetbrains.io.jsonRpc;
import com.intellij.openapi.application.ApplicationManager;
import com.intellij.openapi.components.ServiceManager;
import com.intellij.openapi.extensions.AbstractExtensionPointBean;
import com.intellij.openapi.extensions.ExtensionPointName;
import com.intellij.openapi.util.AtomicNotNullLazyValue;
import com.intellij.openapi.util.NotNullLazyValue;
import com.intellij.util.xmlb.annotations.Attribute;
import org.jetbrains.annotations.NotNull;
public class JsonRpcDomainBean extends AbstractExtensionPointBean {
public static final ExtensionPointName<JsonRpcDomainBean> EP_NAME = ExtensionPointName.create("org.jetbrains.jsonRpcDomain");
@Attribute("name")
public String name;
@Attribute("implementation")
public String implementation;
@Attribute("service")
public String service;
@Attribute("asInstance")
public boolean asInstance = true;
@Attribute("overridable")
public boolean overridable;
private NotNullLazyValue<?> value;
@NotNull
public NotNullLazyValue<?> getValue() {
if (value == null) {
value = new AtomicNotNullLazyValue<Object>() {
@NotNull
@Override
protected Object compute() {
try {
if (service == null) {
Class<Object> aClass = findClass(implementation);
return asInstance ? instantiate(aClass, ApplicationManager.getApplication().getPicoContainer(), true) : aClass;
}
else {
return ServiceManager.getService(findClass(service));
}
}
catch (ClassNotFoundException e) {
throw new RuntimeException(e);
}
}
};
}
return value;
}
}
@@ -18,7 +18,8 @@ import gnu.trove.THashMap;
import gnu.trove.TIntArrayList;
import gnu.trove.TIntProcedure;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.buffer.ByteBufAllocator;
import io.netty.buffer.ByteBufUtil;
import io.netty.util.CharsetUtil;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
@@ -30,6 +31,7 @@ import org.jetbrains.io.JsonUtil;
import java.io.IOException;
import java.lang.reflect.Method;
import java.lang.reflect.Type;
import java.nio.CharBuffer;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
@@ -88,13 +90,16 @@ public class JsonRpcServer implements MessageServer {
}
@Override
public void message(@NotNull Client client, @NotNull String message) throws IOException {
public void messageReceived(@NotNull Client client, @NotNull CharSequence message, boolean isBinary) throws IOException {
if (LOG.isDebugEnabled()) {
LOG.debug("IN " + message);
}
JsonReaderEx reader = new JsonReaderEx(message);
reader.beginArray();
if (!isBinary) {
reader.beginArray();
}
int messageId = reader.peek() == JsonToken.NUMBER ? reader.nextInt() : -1;
String domainName = reader.nextString();
if (domainName.length() == 1) {
@@ -121,7 +126,7 @@ public class JsonRpcServer implements MessageServer {
Object domain = domainHolder.getValue();
String command = reader.nextString();
if (domain instanceof JsonServiceInvocator) {
((JsonServiceInvocator)domain).invoke(command, client, reader, message, messageId);
((JsonServiceInvocator)domain).invoke(command, client, reader, messageId);
return;
}
@@ -135,8 +140,9 @@ public class JsonRpcServer implements MessageServer {
parameters = ArrayUtilRt.EMPTY_OBJECT_ARRAY;
}
reader.endArray();
LOG.assertTrue(reader.peek() == JsonToken.END_DOCUMENT);
if (!isBinary) {
reader.endArray();
}
try {
boolean isStatic = domain instanceof Class;
@@ -152,7 +158,7 @@ public class JsonRpcServer implements MessageServer {
method.setAccessible(true);
Object result = method.invoke(isStatic ? null : domain, parameters);
if (messageId != -1) {
client.send(encodeMessage(messageId, null, null, null, new Object[]{result}));
client.send(encodeMessage(client.getByteBufAllocator(), messageId, null, null, null, new Object[]{result}));
}
return;
}
@@ -166,13 +172,14 @@ public class JsonRpcServer implements MessageServer {
}
public void sendResponse(int messageId, @NotNull Client client, @Nullable CharSequence rawMessage) {
client.send(encodeMessage(messageId, null, null, rawMessage, ArrayUtil.EMPTY_OBJECT_ARRAY));
client.send(encodeMessage(client.getByteBufAllocator(), messageId, null, null, rawMessage, ArrayUtil.EMPTY_OBJECT_ARRAY));
}
public void sendErrorResponse(int messageId, @NotNull Client client, @Nullable CharSequence rawMessage) {
client.send(encodeMessage(messageId, "e", null, rawMessage, ArrayUtil.EMPTY_OBJECT_ARRAY));
client.send(encodeMessage(client.getByteBufAllocator(), messageId, "e", null, rawMessage, ArrayUtil.EMPTY_OBJECT_ARRAY));
}
@SuppressWarnings("unused")
public void sendToClients(@NotNull String domain, @NotNull String name) {
sendToClients(domain, name, null);
}
@@ -184,11 +191,11 @@ public class JsonRpcServer implements MessageServer {
}
private <T> void sendToClients(int messageId, @Nullable String domain, @Nullable String command, @Nullable List<AsyncPromise<Pair<Client, T>>> results, Object[] params) {
clientManager.send(messageId, encodeMessage(messageId, domain, command, null, params), results);
clientManager.send(messageId, encodeMessage(ByteBufAllocator.DEFAULT, messageId, domain, command, null, params), results);
}
public boolean sendWithRawPart(@NotNull Client client, @NotNull String domain, @NotNull String command, @Nullable CharSequence rawMessage, Object... params) {
client.send(encodeMessage(-1, domain, command, rawMessage, params));
client.send(encodeMessage(client.getByteBufAllocator(), -1, domain, command, rawMessage, params));
return true;
}
@@ -199,14 +206,15 @@ public class JsonRpcServer implements MessageServer {
@NotNull
public <T> Promise<T> call(@NotNull Client client, @NotNull String domain, @NotNull String command, Object... params) {
int messageId = messageIdCounter.getAndIncrement();
ByteBuf message = encodeMessage(messageId, domain, command, null, params);
ByteBuf message = encodeMessage(client.getByteBufAllocator(), messageId, domain, command, null, params);
AsyncPromise<T> result = client.send(messageId, message);
LOG.assertTrue(result != null);
return result;
}
@NotNull
private ByteBuf encodeMessage(int messageId,
private ByteBuf encodeMessage(@NotNull ByteBufAllocator byteBufAllocator,
int messageId,
@Nullable String domain,
@Nullable String command,
@Nullable CharSequence rawData,
@@ -217,7 +225,7 @@ public class JsonRpcServer implements MessageServer {
if (LOG.isDebugEnabled()) {
LOG.debug("OUT " + sb.toString());
}
return Unpooled.copiedBuffer(sb, CharsetUtil.UTF_8);
return ByteBufUtil.encodeString(byteBufAllocator, CharBuffer.wrap(sb), CharsetUtil.UTF_8);
}
catch (IOException e) {
throw new RuntimeException(e);
@@ -6,5 +6,5 @@ import org.jetbrains.io.JsonReaderEx;
import java.io.IOException;
public interface JsonServiceInvocator {
void invoke(@NotNull String command, @NotNull Client client, @NotNull JsonReaderEx reader, @NotNull String message, int messageId) throws IOException;
void invoke(@NotNull String command, @NotNull Client client, @NotNull JsonReaderEx reader, int messageId) throws IOException;
}
@@ -5,5 +5,5 @@ import org.jetbrains.annotations.NotNull;
import java.io.IOException;
public interface MessageServer {
void message(@NotNull Client client, @NotNull String message) throws IOException;
void messageReceived(@NotNull Client client, @NotNull CharSequence message, boolean isBinary) throws IOException;
}
@@ -0,0 +1,129 @@
package org.jetbrains.io.jsonRpc.socket;
import com.intellij.openapi.Disposable;
import com.intellij.openapi.diagnostic.Logger;
import com.intellij.openapi.util.AtomicNotNullLazyValue;
import com.intellij.openapi.util.Disposer;
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelHandlerContext;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.ide.BinaryRequestHandler;
import org.jetbrains.ide.BuiltInServerManager;
import org.jetbrains.io.MessageDecoder;
import org.jetbrains.io.jsonRpc.*;
import java.util.List;
import java.util.Map;
import java.util.UUID;
class RpcBinaryRequestHandler extends BinaryRequestHandler implements ExceptionHandler, ClientListener {
private static final Logger LOG = Logger.getInstance(RpcBinaryRequestHandler.class);
private static final UUID ID = UUID.fromString("69957EEB-AFB8-4036-A9A8-00D2D022F9BD");
private final AtomicNotNullLazyValue<ClientManager> clientManager = new AtomicNotNullLazyValue<ClientManager>() {
@NotNull
@Override
protected ClientManager compute() {
ClientManager result = new ClientManager(RpcBinaryRequestHandler.this, RpcBinaryRequestHandler.this);
Disposable serverDisposable = BuiltInServerManager.getInstance().getServerDisposable();
assert serverDisposable != null;
Disposer.register(serverDisposable, result);
rpcServer = new JsonRpcServer(result);
return result;
}
};
private JsonRpcServer rpcServer;
@NotNull
@Override
public UUID getId() {
return ID;
}
@NotNull
@Override
public ChannelHandler getInboundHandler(@NotNull ChannelHandlerContext context) {
SocketClient client = new SocketClient(context.channel());
context.attr(ClientManager.CLIENT).set(client);
clientManager.getValue().addClient(client);
connected(client, null);
return new MyDecoder(client);
}
@Override
public void exceptionCaught(@NotNull Throwable e) {
LOG.error(e);
}
@Override
public void connected(@NotNull Client client, @Nullable Map<String, List<String>> parameters) {
}
@Override
public void disconnected(@NotNull Client client) {
}
private enum State {
LENGTH, CONTENT
}
private class MyDecoder extends MessageDecoder {
private State state = State.LENGTH;
private int contentLength;
private final SocketClient client;
public MyDecoder(@NotNull SocketClient client) {
this.client = client;
}
@Override
protected void messageReceived(@NotNull ChannelHandlerContext context, @NotNull ByteBuf input) throws Exception {
while (true) {
switch (state) {
case LENGTH: {
ByteBuf buffer = getBufferIfSufficient(input, 4, context);
if (buffer == null) {
return;
}
state = State.CONTENT;
contentLength = buffer.readInt();
}
case CONTENT: {
CharSequence content = readChars(input);
if (content == null) {
return;
}
try {
rpcServer.messageReceived(client, content, true);
}
catch (Throwable e) {
clientManager.getValue().exceptionHandler.exceptionCaught(e);
}
finally {
contentLength = 0;
state = State.LENGTH;
}
}
}
}
}
@Override
public void channelInactive(ChannelHandlerContext context) throws Exception {
Client client = context.attr(ClientManager.CLIENT).get();
// if null, so, has already been explicitly removed
if (client != null) {
clientManager.getValue().disconnectClient(context, client, false);
}
}
}
}
@@ -0,0 +1,32 @@
package org.jetbrains.io.jsonRpc.socket;
import io.netty.buffer.ByteBuf;
import io.netty.channel.Channel;
import io.netty.channel.ChannelFuture;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.io.jsonRpc.Client;
import java.nio.channels.ClosedChannelException;
public class SocketClient extends Client {
protected SocketClient(@NotNull Channel channel) {
super(channel);
}
@Override
public ChannelFuture send(@NotNull ByteBuf message) {
if (channel.isOpen()) {
ByteBuf lengthBuffer = channel.alloc().buffer(4);
lengthBuffer.writeInt(message.readableBytes());
channel.write(lengthBuffer);
return channel.writeAndFlush(message);
}
else {
return channel.newFailedFuture(new ClosedChannelException());
}
}
@Override
public void sendHeartbeat() {
}
}
@@ -10,8 +10,6 @@ import org.jetbrains.io.jsonRpc.Client;
import org.jetbrains.io.jsonRpc.ClientManager;
import org.jetbrains.io.jsonRpc.MessageServer;
import java.io.IOException;
@ChannelHandler.Sharable
final class MessageChannelHandler extends SimpleChannelInboundHandlerAdapter<WebSocketFrame> {
private final ClientManager clientManager;
@@ -40,12 +38,11 @@ final class MessageChannelHandler extends SimpleChannelInboundHandlerAdapter<Web
context.channel().writeAndFlush(new PongWebSocketFrame(message.content()));
}
else if (message instanceof TextWebSocketFrame) {
String text = ChannelBufferToString.readString(message.content());
try {
messageServer.message(client, text);
messageServer.messageReceived(client, ChannelBufferToString.readChars(message.content()), false);
}
catch (Throwable e) {
clientManager.exceptionHandler.exceptionCaught(new IOException("Exception while handle message: " + text, e));
clientManager.exceptionHandler.exceptionCaught(e);
}
}
else if (!(message instanceof PongWebSocketFrame)) {
@@ -23,10 +23,12 @@ class WebSocketClient extends Client {
@Override
public ChannelFuture send(@NotNull ByteBuf message) {
if (!channel.isOpen()) {
if (channel.isOpen()) {
return channel.writeAndFlush(new TextWebSocketFrame(message));
}
else {
return channel.newFailedFuture(new ClosedChannelException());
}
return channel.writeAndFlush(new TextWebSocketFrame(message));
}
@Override
@@ -8,13 +8,14 @@ import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelFutureListener;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.http.FullHttpRequest;
import io.netty.handler.codec.http.HttpHeaders;
import io.netty.handler.codec.http.HttpHeaderNames;
import io.netty.handler.codec.http.HttpMethod;
import io.netty.handler.codec.http.QueryStringDecoder;
import io.netty.handler.codec.http.websocketx.WebSocketFrameAggregator;
import io.netty.handler.codec.http.websocketx.WebSocketServerHandshaker;
import io.netty.handler.codec.http.websocketx.WebSocketServerHandshakerFactory;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.ide.BuiltInServerManager;
import org.jetbrains.ide.HttpRequestHandler;
import org.jetbrains.io.BuiltInServer;
@@ -27,7 +28,7 @@ import java.util.Map;
public abstract class WebSocketHandshakeHandler extends HttpRequestHandler implements ClientListener, ExceptionHandler {
private static final Logger LOG = Logger.getInstance(WebSocketHandshakeHandler.class);
private final AtomicNotNullLazyValue<ClientManager> server = new AtomicNotNullLazyValue<ClientManager>() {
private final AtomicNotNullLazyValue<ClientManager> clientManager = new AtomicNotNullLazyValue<ClientManager>() {
@NotNull
@Override
protected ClientManager compute() {
@@ -43,7 +44,7 @@ public abstract class WebSocketHandshakeHandler extends HttpRequestHandler imple
@Override
public boolean isSupported(@NotNull FullHttpRequest request) {
return request.method() == HttpMethod.GET &&
"WebSocket".equalsIgnoreCase(request.headers().get(HttpHeaders.Names.UPGRADE)) &&
"WebSocket".equalsIgnoreCase(request.headers().get(HttpHeaderNames.UPGRADE)) &&
request.uri().length() > 2;
}
@@ -51,7 +52,7 @@ public abstract class WebSocketHandshakeHandler extends HttpRequestHandler imple
}
@Override
public void exceptionCaught(Throwable e) {
public void exceptionCaught(@NotNull Throwable e) {
NettyUtil.log(e, LOG);
}
@@ -61,10 +62,8 @@ public abstract class WebSocketHandshakeHandler extends HttpRequestHandler imple
return true;
}
private void handleWebSocketRequest(final ChannelHandlerContext context,
FullHttpRequest request,
final QueryStringDecoder uriDecoder) {
WebSocketServerHandshakerFactory factory = new WebSocketServerHandshakerFactory("ws://" + HttpHeaders.getHost(request) + uriDecoder.path(), null, false, NettyUtil.MAX_CONTENT_LENGTH);
private void handleWebSocketRequest(@NotNull final ChannelHandlerContext context, @NotNull FullHttpRequest request, @NotNull final QueryStringDecoder uriDecoder) {
WebSocketServerHandshakerFactory factory = new WebSocketServerHandshakerFactory("ws://" + request.headers().get(HttpHeaderNames.HOST) + uriDecoder.path(), null, false, NettyUtil.MAX_CONTENT_LENGTH);
WebSocketServerHandshaker handshaker = factory.newHandshaker(request);
if (handshaker == null) {
WebSocketServerHandshakerFactory.sendUnsupportedVersionResponse(context.channel());
@@ -81,9 +80,9 @@ public abstract class WebSocketHandshakeHandler extends HttpRequestHandler imple
@Override
public void operationComplete(ChannelFuture future) throws Exception {
if (future.isSuccess()) {
ClientManager webSocketServer = server.getValue();
webSocketServer.addClient(client);
MessageChannelHandler messageChannelHandler = new MessageChannelHandler(webSocketServer, getMessageServer());
ClientManager clientManager = WebSocketHandshakeHandler.this.clientManager.getValue();
clientManager.addClient(client);
MessageChannelHandler messageChannelHandler = new MessageChannelHandler(clientManager, getMessageServer());
BuiltInServer.replaceDefaultHandler(context, messageChannelHandler);
ChannelHandlerContext messageChannelHandlerContext = context.pipeline().context(messageChannelHandler);
context.pipeline().addBefore(messageChannelHandlerContext.name(), "webSocketFrameAggregator", new WebSocketFrameAggregator(NettyUtil.MAX_CONTENT_LENGTH));
@@ -98,6 +97,6 @@ public abstract class WebSocketHandshakeHandler extends HttpRequestHandler imple
protected abstract MessageServer getMessageServer();
@Override
public void connected(@NotNull Client client, Map<String, List<String>> parameters) {
public void connected(@NotNull Client client, @Nullable Map<String, List<String>> parameters) {
}
}
@@ -18,10 +18,6 @@ public final class SingletonNotificationManager {
private Runnable expiredListener;
public SingletonNotificationManager(@NotNull String groupId, @NotNull NotificationType type, @Nullable NotificationListener listener) {
this(new NotificationGroup(groupId, NotificationDisplayType.STICKY_BALLOON, true), type, listener);
}
public SingletonNotificationManager(@NotNull NotificationGroup group, @NotNull NotificationType type, @Nullable NotificationListener listener) {
this.group = group;
this.type = type;
@@ -18,6 +18,7 @@ import junit.framework.TestCase
import org.junit.rules.RuleChain
import org.junit.Rule
import org.junit.Test
import org.jetbrains.io.MessageDecoder
public class BinaryRequestHandlerTest {
private val fixtureManager = FixtureRule()
@@ -36,16 +37,10 @@ public class BinaryRequestHandlerTest {
val bootstrap = NettyUtil.oioClientBootstrap().handler(object : ChannelInitializer<Channel>() {
override fun initChannel(channel: Channel) {
channel.pipeline().addLast(object : Decoder() {
override fun messageReceived(context: ChannelHandlerContext, message: ByteBuf) {
override fun messageReceived(context: ChannelHandlerContext, input: ByteBuf) {
val requiredLength = 4 + text.length()
val buffer = getBufferIfSufficient(message, requiredLength, context)
if (buffer == null) {
message.release()
}
else {
val response = buffer.toString(buffer.readerIndex(), requiredLength, CharsetUtil.UTF_8)
buffer.skipBytes(requiredLength)
buffer.release()
val response = readContent(input, context, requiredLength) {(buffer, context) -> buffer.toString(buffer.readerIndex(), requiredLength, CharsetUtil.UTF_8) }
if (response != null) {
result.setDone(response)
}
}
@@ -90,55 +85,37 @@ public class BinaryRequestHandlerTest {
return ID
}
override fun getInboundHandler(): ChannelHandler {
override fun getInboundHandler(context: ChannelHandlerContext): ChannelHandler {
return MyDecoder()
}
private class MyDecoder : Decoder() {
private class MyDecoder : MessageDecoder() {
private var state = State.HEADER
private var contentLength = -1
private enum class State {
HEADER
CONTENT
}
override fun messageReceived(context: ChannelHandlerContext, message: ByteBuf) {
override fun messageReceived(context: ChannelHandlerContext, input: ByteBuf) {
while (true) {
when (state) {
State.HEADER -> {
run {
val buffer = getBufferIfSufficient(message, 2, context)
if (buffer == null) {
message.release()
return
}
contentLength = buffer.readUnsignedShort()
state = State.CONTENT
val buffer = getBufferIfSufficient(input, 2, context)
if (buffer == null) {
return
}
run {
val buffer = getBufferIfSufficient(message, contentLength, context)
if (buffer == null) {
message.release()
return
}
val messageText = buffer.toString(buffer.readerIndex(), contentLength, CharsetUtil.UTF_8)
buffer.skipBytes(contentLength)
state = State.HEADER
context.writeAndFlush(Unpooled.copiedBuffer("got-" + messageText, CharsetUtil.UTF_8))
}
contentLength = buffer.readUnsignedShort()
state = State.CONTENT
}
State.CONTENT -> {
val buffer = getBufferIfSufficient(message, contentLength, context)
if (buffer == null) {
message.release()
val messageText = readChars(input)
if (messageText == null) {
return
}
val messageText = buffer.toString(buffer.readerIndex(), contentLength, CharsetUtil.UTF_8)
buffer.skipBytes(contentLength)
state = State.HEADER
context.writeAndFlush(Unpooled.copiedBuffer("got-" + messageText, CharsetUtil.UTF_8))
}
@@ -5,7 +5,7 @@
<!-- otherwise plugin will be not loaded -->
<depends>com.intellij.modules.xml</depends>
<extensions defaultExtensionNs="com.intellij">
<extensions defaultExtensionNs="org.jetbrains">
<binaryRequestHandler implementation="org.jetbrains.ide.BinaryRequestHandlerTest$MyBinaryRequestHandler"/>
</extensions>
</idea-plugin>
@@ -1,143 +0,0 @@
package com.intellij.openapi.util;
import com.intellij.openapi.Disposable;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.io.IOException;
import java.util.concurrent.atomic.AtomicReference;
public abstract class AsyncValueLoader<T> {
private final AtomicReference<AsyncResult<T>> ref = new AtomicReference<AsyncResult<T>>();
private volatile long modificationCount;
private volatile long loadedModificationCount;
private final Runnable doneHandler = new Runnable() {
@Override
public void run() {
loadedModificationCount = modificationCount;
}
};
@NotNull
public final AsyncResult<T> get() {
return get(true);
}
public final void reset() {
AsyncResult<T> oldValue = ref.getAndSet(null);
if (oldValue != null) {
rejectAndDispose(oldValue);
}
}
private void rejectAndDispose(@NotNull AsyncResult<T> asyncResult) {
try {
if (!asyncResult.isProcessed()) {
asyncResult.setRejected();
}
}
finally {
T result = asyncResult.getResult();
if (result != null) {
disposeResult(result);
}
}
}
protected void disposeResult(@NotNull T result) {
if (result instanceof Disposable) {
Disposer.dispose((Disposable)result, false);
}
}
public final boolean has() {
AsyncResult<T> result = ref.get();
return result != null && result.isDone() && result.getResult() != null;
}
@NotNull
public final AsyncResult<T> get(boolean checkFreshness) {
AsyncResult<T> asyncResult = ref.get();
if (asyncResult == null) {
if (!ref.compareAndSet(null, asyncResult = new AsyncResult<T>())) {
return ref.get();
}
}
else if (!asyncResult.isProcessed()) {
// if current asyncResult is not processed, so, we don't need to check cache state
return asyncResult;
}
else if (asyncResult.isDone()) {
if (!checkFreshness || isUpToDate(asyncResult.getResult())) {
return asyncResult;
}
if (!ref.compareAndSet(asyncResult, asyncResult = new AsyncResult<T>())) {
AsyncResult<T> valueFromAnotherThread = ref.get();
while (valueFromAnotherThread == null) {
if (ref.compareAndSet(null, asyncResult)) {
callLoad(asyncResult);
return asyncResult;
}
else {
valueFromAnotherThread = ref.get();
}
}
return valueFromAnotherThread;
}
}
callLoad(asyncResult);
return asyncResult;
}
/**
* if result was rejected, by default this result will not be canceled - call get() will return rejected result instead of attempt to load again,
* but you can change this behavior - return true if you want to cancel result on reject
*/
protected boolean isCancelOnReject() {
return false;
}
private void callLoad(final @NotNull AsyncResult<T> result) {
if (isCancelOnReject()) {
result.doWhenRejected(new Runnable() {
@Override
public void run() {
ref.compareAndSet(result, null);
}
});
}
result.doWhenDone(doneHandler);
try {
load(result);
}
catch (Throwable e) {
ref.compareAndSet(result, null);
rejectAndDispose(result);
//noinspection InstanceofCatchParameter
throw e instanceof RuntimeException ? ((RuntimeException)e) : new RuntimeException(e);
}
}
protected abstract void load(@NotNull AsyncResult<T> result) throws IOException;
protected boolean isUpToDate(@Nullable T result) {
return loadedModificationCount == modificationCount;
}
public final void set(@NotNull T result) {
AsyncResult<T> oldValue = ref.getAndSet(AsyncResult.done(result));
if (oldValue != null) {
rejectAndDispose(oldValue);
}
}
public final void markDirty() {
modificationCount++;
}
}
@@ -2,9 +2,10 @@ package io.netty.bootstrap;
import io.netty.channel.Channel;
import io.netty.channel.ChannelFuture;
import org.jetbrains.annotations.NotNull;
public final class BootstrapUtil {
public static ChannelFuture initAndRegister(Channel channel, Bootstrap bootstrap) throws Throwable {
public static ChannelFuture initAndRegister(@NotNull Channel channel, @NotNull Bootstrap bootstrap) throws Throwable {
try {
bootstrap.init(channel);
}
@@ -13,9 +14,9 @@ public final class BootstrapUtil {
throw e;
}
ChannelFuture regFuture = bootstrap.group().register(channel);
ChannelFuture registrationFuture = bootstrap.group().register(channel);
//noinspection ThrowableResultOfMethodCallIgnored
if (regFuture.cause() != null) {
if (registrationFuture.cause() != null) {
if (channel.isRegistered()) {
channel.close();
}
@@ -23,6 +24,6 @@ public final class BootstrapUtil {
channel.unsafe().closeForcibly();
}
}
return regFuture;
return registrationFuture;
}
}
@@ -152,6 +152,8 @@ public class AsyncPromise<T> extends Promise<T> implements Getter<T> {
@Override
void notify(@NotNull final AsyncPromise<T> child) {
LOG.assertTrue(child != this);
switch (state) {
case PENDING:
break;
@@ -300,9 +302,9 @@ public class AsyncPromise<T> extends Promise<T> implements Getter<T> {
return done instanceof Obsolescent && ((Obsolescent)done).isObsolete();
}
public void setError(@NotNull Throwable error) {
public boolean setError(@NotNull Throwable error) {
if (state != State.PENDING) {
return;
return false;
}
result = error;
@@ -318,6 +320,7 @@ public class AsyncPromise<T> extends Promise<T> implements Getter<T> {
else if (!(error instanceof MessageError)) {
LOG.error(error);
}
return true;
}
private void clearHandlers() {
@@ -326,7 +329,7 @@ public class AsyncPromise<T> extends Promise<T> implements Getter<T> {
}
@Override
public void processed(@NotNull final Consumer<T> processed) {
public Promise<T> processed(@NotNull final Consumer<T> processed) {
done(processed);
rejected(new Consumer<Throwable>() {
@Override
@@ -334,5 +337,6 @@ public class AsyncPromise<T> extends Promise<T> implements Getter<T> {
processed.consume(null);
}
});
return this;
}
}
@@ -0,0 +1,162 @@
package org.jetbrains.concurrency;
import com.intellij.openapi.Disposable;
import com.intellij.openapi.util.Disposer;
import com.intellij.openapi.util.Getter;
import com.intellij.util.Consumer;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.io.IOException;
import java.util.concurrent.atomic.AtomicReference;
public abstract class AsyncValueLoader<T> {
private final AtomicReference<Promise<T>> ref = new AtomicReference<Promise<T>>();
private volatile long modificationCount;
private volatile long loadedModificationCount;
private final Consumer<T> doneHandler = new Consumer<T>() {
@Override
public void consume(T o) {
loadedModificationCount = modificationCount;
}
};
@NotNull
public final Promise<T> get() {
return get(true);
}
public final T getResult() {
//noinspection unchecked
return ((Getter<T>)get(true)).get();
}
public final void reset() {
Promise<T> oldValue = ref.getAndSet(null);
if (oldValue instanceof AsyncPromise) {
rejectAndDispose((AsyncPromise<T>)oldValue);
}
}
private void rejectAndDispose(@NotNull AsyncPromise<T> asyncResult) {
try {
asyncResult.setError(Promise.createError("rejected"));
}
finally {
T result = asyncResult.get();
if (result != null) {
disposeResult(result);
}
}
}
protected void disposeResult(@NotNull T result) {
if (result instanceof Disposable) {
Disposer.dispose((Disposable)result, false);
}
}
public final boolean has() {
Promise<T> result = ref.get();
//noinspection unchecked
return result != null && result.getState() == Promise.State.FULFILLED && ((Getter<T>)result).get() != null;
}
@NotNull
public final Promise<T> get(boolean checkFreshness) {
Promise<T> promise = ref.get();
if (promise == null) {
if (!ref.compareAndSet(null, promise = new AsyncPromise<T>())) {
return ref.get();
}
}
else {
Promise.State state = promise.getState();
if (state == Promise.State.PENDING) {
// if current promise is not processed, so, we don't need to check cache state
return promise;
}
else if (state == Promise.State.FULFILLED) {
//noinspection unchecked
if (!checkFreshness || isUpToDate(((Getter<T>)promise).get())) {
return promise;
}
if (!ref.compareAndSet(promise, promise = new AsyncPromise<T>())) {
Promise<T> valueFromAnotherThread = ref.get();
while (valueFromAnotherThread == null) {
if (ref.compareAndSet(null, promise)) {
return getPromise((AsyncPromise<T>)promise);
}
else {
valueFromAnotherThread = ref.get();
}
}
return valueFromAnotherThread;
}
}
}
return getPromise((AsyncPromise<T>)promise);
}
/**
* if result was rejected, by default this result will not be canceled - call get() will return rejected result instead of attempt to load again,
* but you can change this behavior - return true if you want to cancel result on reject
*/
protected boolean isCancelOnReject() {
return false;
}
@NotNull
private Promise<T> getPromise(@NotNull AsyncPromise<T> promise) {
final Promise<T> effectivePromise;
try {
effectivePromise = load(promise);
if (effectivePromise != promise) {
ref.compareAndSet(promise, effectivePromise);
}
}
catch (Throwable e) {
ref.compareAndSet(promise, null);
rejectAndDispose(promise);
//noinspection InstanceofCatchParameter
throw e instanceof RuntimeException ? ((RuntimeException)e) : new RuntimeException(e);
}
effectivePromise.done(doneHandler);
if (isCancelOnReject()) {
effectivePromise.rejected(new Consumer<Throwable>() {
@Override
public void consume(Throwable throwable) {
ref.compareAndSet(effectivePromise, null);
}
});
}
if (effectivePromise != promise) {
effectivePromise.notify(promise);
}
return effectivePromise;
}
@NotNull
protected abstract Promise<T> load(@NotNull AsyncPromise<T> result) throws IOException;
protected boolean isUpToDate(@Nullable T result) {
return loadedModificationCount == modificationCount;
}
public final void set(@NotNull T result) {
Promise<T> oldValue = ref.getAndSet(Promise.resolve(result));
if (oldValue != null && oldValue instanceof AsyncPromise) {
rejectAndDispose((AsyncPromise<T>)oldValue);
}
}
public final void markDirty() {
modificationCount++;
}
}
@@ -1,17 +0,0 @@
package org.jetbrains.concurrency;
import com.intellij.util.Consumer;
import com.intellij.util.Function;
public abstract class ConsumerRunnable implements Consumer<Void>, Runnable, Function<Object, Void> {
@Override
public final void consume(Void v) {
run();
}
@Override
public final Void fun(Object object) {
run();
return null;
}
}
@@ -29,8 +29,9 @@ class DonePromise<T> extends Promise<T> implements Getter<T> {
}
@Override
public void processed(@NotNull Consumer<T> processed) {
public Promise<T> processed(@NotNull Consumer<T> processed) {
done(processed);
return this;
}
@NotNull
@@ -117,25 +117,13 @@ public abstract class Promise<T> {
@NotNull
public abstract Promise<T> done(@NotNull Consumer<T> done);
@NotNull
public Promise<T> done(@NotNull ConsumerRunnable done) {
//noinspection unchecked
return done((Consumer<T>)done);
}
@NotNull
public abstract Promise<T> processed(@NotNull final AsyncPromise<T> fulfilled);
@NotNull
public Promise<Void> then(@NotNull ConsumerRunnable done) {
//noinspection unchecked
return then((Function<T, Void>)done);
}
@NotNull
public abstract Promise<T> rejected(@NotNull Consumer<Throwable> rejected);
public abstract void processed(@NotNull Consumer<T> processed);
public abstract Promise<T> processed(@NotNull Consumer<T> processed);
@NotNull
public abstract <SUB_RESULT> Promise<SUB_RESULT> then(@NotNull Function<T, SUB_RESULT> done);
@@ -22,6 +22,7 @@ public abstract class PromiseManager<HOST, VALUE> {
return true;
}
@NotNull
public abstract Promise<VALUE> load(@NotNull HOST host);
public final void reset(HOST host) {
@@ -103,8 +104,10 @@ public abstract class PromiseManager<HOST, VALUE> {
}
Promise<VALUE> effectivePromise = load(host);
fieldUpdater.compareAndSet(host, promise, effectivePromise);
effectivePromise.notify((AsyncPromise<VALUE>)promise);
if (effectivePromise != promise) {
fieldUpdater.compareAndSet(host, promise, effectivePromise);
effectivePromise.notify((AsyncPromise<VALUE>)promise);
}
return effectivePromise;
}
}
@@ -34,8 +34,9 @@ class RejectedPromise<T> extends Promise<T> {
}
@Override
public void processed(@NotNull Consumer<T> processed) {
public RejectedPromise<T> processed(@NotNull Consumer<T> processed) {
processed.consume(null);
return this;
}
@NotNull
@@ -1,6 +1,6 @@
package org.jetbrains.io;
import com.intellij.util.text.StringFactory;
import com.intellij.util.text.CharArrayCharSequence;
import io.netty.buffer.ByteBuf;
import io.netty.util.CharsetUtil;
import org.jetbrains.annotations.NotNull;
@@ -13,42 +13,31 @@ import java.nio.charset.CharsetDecoder;
import java.nio.charset.CoderResult;
public final class ChannelBufferToString {
public static String readString(ByteBuf buffer) {
return charBufferToString(readIntoCharBuffer(null, buffer, buffer.readableBytes()));
@NotNull
public static CharSequence readChars(@NotNull ByteBuf buffer) throws CharacterCodingException {
return new MyCharArrayCharSequence(readIntoCharBuffer(CharsetUtil.getDecoder(CharsetUtil.UTF_8), buffer, buffer.readableBytes(), null));
}
public static String readString(ByteBuf buffer, int byteCount) {
return charBufferToString(readIntoCharBuffer(null, buffer, byteCount));
@SuppressWarnings("unused")
@NotNull
public static CharSequence readChars(@NotNull ByteBuf buffer, int byteCount) throws CharacterCodingException {
return new MyCharArrayCharSequence(readIntoCharBuffer(CharsetUtil.getDecoder(CharsetUtil.UTF_8), buffer, byteCount, null));
}
public static String charBufferToString(CharBuffer charBuffer) {
char[] array = charBuffer.array();
if (array.length == charBuffer.position()) {
return StringFactory.createShared(array);
}
else {
return charBuffer.flip().toString();
}
}
public static CharBuffer readIntoCharBuffer(@Nullable CharBuffer charBuffer, @NotNull ByteBuf buffer, int byteCount) {
CharsetDecoder decoder = CharsetUtil.getDecoder(CharsetUtil.UTF_8);
@NotNull
public static CharBuffer readIntoCharBuffer(@NotNull CharsetDecoder decoder, @NotNull ByteBuf buffer, int byteCount, @Nullable CharBuffer charBuffer) throws CharacterCodingException {
ByteBuffer in = buffer.nioBuffer(buffer.readerIndex(), byteCount);
if (charBuffer == null) {
charBuffer = CharBuffer.allocate((int) ((double) in.remaining() * decoder.maxCharsPerByte()));
charBuffer = CharBuffer.allocate((int)((float)in.remaining() * decoder.maxCharsPerByte()));
}
try {
CoderResult cr = decoder.decode(in, charBuffer, true);
if (!cr.isUnderflow()) {
cr.throwException();
}
cr = decoder.flush(charBuffer);
if (!cr.isUnderflow()) {
cr.throwException();
}
CoderResult cr = decoder.decode(in, charBuffer, true);
if (!cr.isUnderflow()) {
cr.throwException();
}
catch (CharacterCodingException x) {
throw new IllegalStateException(x);
cr = decoder.flush(charBuffer);
if (!cr.isUnderflow()) {
cr.throwException();
}
buffer.skipBytes(byteCount);
@@ -61,4 +50,18 @@ public final class ChannelBufferToString {
buffer.writeByte(string.charAt(i));
}
}
// we can produce char sequence CharSequence result = CharsetUtil.UTF_8.decode(buffer.nioBuffer(buffer.readerIndex(), required));
// but later, in JsonReaderEx, it will be toString in any case, so, in this case, intermediate java.nio.HeapCharBuffer will be created - so, we stay with String
// we must return string on subSequence() - JsonReaderEx will call toString in any case
public static final class MyCharArrayCharSequence extends CharArrayCharSequence {
public MyCharArrayCharSequence(@NotNull CharBuffer charBuffer) {
super(charBuffer.array(), charBuffer.arrayOffset(), charBuffer.position());
}
@Override
public CharSequence subSequence(int start, int end) {
return start == 0 && end == length() ? this : new String(myChars, myStart + start, end - start);
}
}
}
@@ -21,8 +21,6 @@ import io.netty.channel.ChannelHandlerAdapter;
import io.netty.channel.ChannelHandlerContext;
import org.jetbrains.annotations.NotNull;
import java.net.ConnectException;
@ChannelHandler.Sharable
public final class ChannelExceptionHandler extends ChannelHandlerAdapter {
private static final Logger LOG = Logger.getInstance(ChannelExceptionHandler.class);
@@ -39,13 +37,6 @@ public final class ChannelExceptionHandler extends ChannelHandlerAdapter {
@Override
public void exceptionCaught(ChannelHandlerContext context, Throwable cause) throws Exception {
// don't report about errors while connecting
// WEB-7727
if (cause instanceof ConnectException) {
LOG.debug(cause);
}
else {
NettyUtil.log(cause, LOG);
}
NettyUtil.logAndClose(cause, LOG, context.channel());
}
}
@@ -0,0 +1,149 @@
/*
* Copyright 2000-2015 JetBrains s.r.o.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.jetbrains.io;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.CompositeByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.io.IOException;
public abstract class Decoder extends ChannelInboundHandlerAdapter {
// Netty MessageAggregator default value
protected static final int DEFAULT_MAX_COMPOSITE_BUFFER_COMPONENTS = 1024;
private ByteBuf cumulation;
@Override
public final void channelRead(ChannelHandlerContext context, Object message) throws Exception {
if (message instanceof ByteBuf) {
ByteBuf input = (ByteBuf)message;
try {
messageReceived(context, input);
}
finally {
input.release();
}
}
else {
context.fireChannelRead(message);
}
}
protected abstract void messageReceived(@NotNull ChannelHandlerContext context, @NotNull ByteBuf input) throws Exception;
public interface FullMessageConsumer<T> {
T contentReceived(@NotNull ByteBuf input, @NotNull ChannelHandlerContext context, boolean isCumulateBuffer) throws IOException;
}
@Nullable
protected final <T> T readContent(@NotNull ByteBuf input, @NotNull ChannelHandlerContext context, int contentLength, @NotNull FullMessageConsumer<T> fullMessageConsumer) throws IOException {
ByteBuf buffer = getBufferIfSufficient(input, contentLength, context);
if (buffer == null) {
return null;
}
boolean isCumulateBuffer = buffer != input;
int oldReaderIndex = input.readerIndex();
try {
return fullMessageConsumer.contentReceived(buffer, context, isCumulateBuffer);
}
finally {
if (isCumulateBuffer) {
// cumulation buffer - release it
buffer.release();
}
else {
buffer.readerIndex(oldReaderIndex + contentLength);
}
}
}
@Nullable
protected final ByteBuf getBufferIfSufficient(@NotNull ByteBuf input, int requiredLength, @NotNull ChannelHandlerContext context) {
if (!input.isReadable()) {
return null;
}
if (cumulation == null) {
if (input.readableBytes() < requiredLength) {
cumulation = input;
input.retain();
input.touch();
return null;
}
else {
return input;
}
}
else {
int currentAccumulatedByteCount = cumulation.readableBytes();
if ((currentAccumulatedByteCount + input.readableBytes()) < requiredLength) {
CompositeByteBuf compositeByteBuf;
if ((cumulation instanceof CompositeByteBuf)) {
compositeByteBuf = (CompositeByteBuf)cumulation;
}
else {
compositeByteBuf = context.alloc().compositeBuffer(DEFAULT_MAX_COMPOSITE_BUFFER_COMPONENTS);
compositeByteBuf.addComponent(cumulation);
cumulation = compositeByteBuf;
}
compositeByteBuf.addComponent(input);
input.retain();
input.touch();
return null;
}
else {
CompositeByteBuf buffer;
if (cumulation instanceof CompositeByteBuf) {
buffer = (CompositeByteBuf)cumulation;
buffer.addComponent(input);
}
else {
// may be it will be used by client to cumulate something - don't set artificial restriction (2)
buffer = context.alloc().compositeBuffer(DEFAULT_MAX_COMPOSITE_BUFFER_COMPONENTS);
buffer.addComponents(cumulation, input);
}
// we don't set writerIndex on addComponent, it is clear to set it to requiredLength here
buffer.writerIndex(requiredLength);
input.skipBytes(requiredLength - currentAccumulatedByteCount);
input.retain();
input.touch();
cumulation = null;
return buffer;
}
}
}
@Override
public void channelInactive(ChannelHandlerContext context) throws Exception {
try {
if (cumulation != null) {
cumulation.release();
cumulation = null;
}
}
finally {
super.channelInactive(context);
}
}
}
@@ -1,53 +1,95 @@
package org.jetbrains.rpc;
package org.jetbrains.io;
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.util.CharsetUtil;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.io.ChannelBufferToString;
import org.jetbrains.io.SimpleChannelInboundHandlerAdapter;
import java.nio.CharBuffer;
import java.nio.charset.CharacterCodingException;
import java.nio.charset.CharsetDecoder;
public abstract class MessageDecoder extends SimpleChannelInboundHandlerAdapter<ByteBuf> {
public abstract class MessageDecoder extends Decoder {
protected int contentLength;
protected final StringBuilder builder = new StringBuilder(64);
private CharBuffer chunkedContent;
private int consumedContentByteCount = 0;
private final CharsetDecoder charsetDecoder = CharsetUtil.getDecoder(CharsetUtil.UTF_8);
protected final int parseContentLength() {
return parseInt(builder, 0, false, 10);
}
@Nullable
protected String doReadContent(@NotNull ByteBuf buffer) {
protected final CharSequence readChars(@NotNull ByteBuf input) throws CharacterCodingException {
if (!input.isReadable()) {
return null;
}
int required = contentLength - consumedContentByteCount;
String result;
if (buffer.readableBytes() < required) {
if (input.readableBytes() < required) {
if (chunkedContent == null) {
chunkedContent = CharBuffer.allocate(contentLength);
chunkedContent = CharBuffer.allocate((int)((float)contentLength * charsetDecoder.maxCharsPerByte()));
}
int count = buffer.readableBytes();
ChannelBufferToString.readIntoCharBuffer(chunkedContent, buffer, count);
int count = input.readableBytes();
ChannelBufferToString.readIntoCharBuffer(charsetDecoder, input, count, chunkedContent);
consumedContentByteCount += count;
return null;
}
else if (chunkedContent != null) {
ChannelBufferToString.readIntoCharBuffer(chunkedContent, buffer, required);
result = ChannelBufferToString.charBufferToString(chunkedContent);
chunkedContent = null;
consumedContentByteCount = 0;
return result;
}
else {
// we can produce char sequence CharSequence result = CharsetUtil.UTF_8.decode(buffer.toByteBuffer(buffer.readerIndex(), required));
// but later, in JsonReaderEx, it will be toString in any case, so, in this case, intermediate java.nio.HeapCharBuffer will be created - so, we stay with String
return ChannelBufferToString.readString(buffer, required);
CharBuffer charBuffer = chunkedContent;
if (charBuffer != null) {
chunkedContent = null;
consumedContentByteCount = 0;
}
return new ChannelBufferToString.MyCharArrayCharSequence(ChannelBufferToString.readIntoCharBuffer(charsetDecoder, input, required, charBuffer));
}
}
@Override
public void channelInactive(ChannelHandlerContext context) throws Exception {
try {
chunkedContent = null;
}
finally {
super.channelInactive(context);
}
}
public static boolean readUntil(char what, @NotNull ByteBuf buffer, @NotNull StringBuilder builder) {
int i = buffer.readerIndex();
//noinspection ForLoopThatDoesntUseLoopVariable
for (int n = buffer.writerIndex(); i < n; i++) {
char c = (char)buffer.getByte(i);
if (c == what) {
buffer.readerIndex(i + 1);
return true;
}
else {
builder.append(c);
}
}
buffer.readerIndex(i);
return false;
}
public static void skipWhitespace(@NotNull ByteBuf buffer) {
int i = buffer.readerIndex();
int n = buffer.writerIndex();
for (; i < n; i++) {
char c = (char)buffer.getByte(i);
if (c != ' ') {
buffer.readerIndex(i);
return;
}
}
buffer.readerIndex(n);
}
/**
* Javolution - Java(TM) Solution for Real-Time and Embedded Systems
* Copyright (C) 2006 - Javolution (http://javolution.org/)
@@ -56,7 +98,7 @@ public abstract class MessageDecoder extends SimpleChannelInboundHandlerAdapter<
* Permission to use, copy, modify, and distribute this software is
* freely granted, provided that this notice is preserved.
*/
private static int parseInt(final CharSequence value, final int start, final boolean isNegative, final int radix) {
public static int parseInt(@NotNull CharSequence value, int start, boolean isNegative, int radix) {
final int end = value.length();
int result = 0; // Accumulates negatively (avoid MIN_VALUE overflow).
int i = start;
@@ -29,8 +29,7 @@ import io.netty.channel.socket.oio.OioSocketChannel;
import io.netty.handler.codec.http.HttpObjectAggregator;
import io.netty.handler.codec.http.HttpRequestDecoder;
import io.netty.handler.codec.http.HttpResponseEncoder;
import io.netty.handler.codec.http.cors.CorsConfig;
import io.netty.handler.codec.http.cors.CorsHandler;
import io.netty.handler.stream.ChunkedWriteHandler;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.concurrency.AsyncPromise;
@@ -39,6 +38,7 @@ import org.jetbrains.ide.PooledThreadExecutor;
import java.io.IOException;
import java.net.BindException;
import java.net.ConnectException;
import java.net.InetSocketAddress;
import java.net.Socket;
import java.util.Random;
@@ -57,7 +57,24 @@ public final class NettyUtil {
}
}
public static void log(Throwable throwable, Logger log) {
public static void logAndClose(@NotNull Throwable error, @NotNull Logger log, @NotNull Channel channel) {
// don't report about errors while connecting
// WEB-7727
try {
if (error instanceof ConnectException) {
log.debug(error);
}
else {
log(error, log);
}
}
finally {
log.info("Channel will be closed due to error");
channel.close();
}
}
public static void log(@NotNull Throwable throwable, @NotNull Logger log) {
if (isAsWarning(throwable)) {
log.warn(throwable);
}
@@ -85,8 +102,7 @@ public final class NettyUtil {
}
@Nullable
private static Channel doConnect(@NotNull Bootstrap bootstrap, @NotNull InetSocketAddress remoteAddress, @Nullable AsyncPromise<?> promise, int maxAttemptCount)
throws Throwable {
private static Channel doConnect(@NotNull Bootstrap bootstrap, @NotNull InetSocketAddress remoteAddress, @Nullable AsyncPromise<?> promise, int maxAttemptCount) throws Throwable {
int attemptCount = 0;
if (bootstrap.group() instanceof NioEventLoopGroup) {
@@ -141,13 +157,12 @@ public final class NettyUtil {
}
}
}
OioSocketChannel channel = new OioSocketChannel(socket);
BootstrapUtil.initAndRegister(channel, bootstrap).awaitUninterruptibly();
BootstrapUtil.initAndRegister(channel, bootstrap).sync();
return channel;
}
private static boolean isAsWarning(Throwable throwable) {
private static boolean isAsWarning(@NotNull Throwable throwable) {
String message = throwable.getMessage();
if (message == null) {
return false;
@@ -195,9 +210,13 @@ public final class NettyUtil {
}
public static void addHttpServerCodec(@NotNull ChannelPipeline pipeline) {
pipeline.addLast(new HttpRequestDecoder(),
new HttpResponseEncoder(),
new CorsHandler(CorsConfig.withAnyOrigin().allowCredentials().allowNullOrigin().allowedRequestMethods().build()),
new HttpObjectAggregator(MAX_CONTENT_LENGTH));
pipeline.addLast("httpRequestEncoder", new HttpResponseEncoder());
pipeline.addLast("httpRequestDecoder", new HttpRequestDecoder());
pipeline.addLast("httpObjectAggregator", new HttpObjectAggregator(MAX_CONTENT_LENGTH));
// could be added earlier if HTTPS
if (pipeline.get(ChunkedWriteHandler.class) == null) {
pipeline.addLast("chunkedWriteHandler", new ChunkedWriteHandler());
}
//pipeline.addLast("corsHandler", new CorsHandler(CorsConfig.withAnyOrigin().allowCredentials().allowNullOrigin().allowedRequestMethods().build()));
}
}
@@ -5,8 +5,12 @@
<extensionPoint qualifiedName="org.jetbrains.webServerRootsProvider" interface="org.jetbrains.builtInWebServer.WebServerRootsProvider"/>
<extensionPoint name="httpRequestHandler" interface="org.jetbrains.ide.HttpRequestHandler"/>
<extensionPoint name="binaryRequestHandler" interface="org.jetbrains.ide.BinaryRequestHandler"/>
<extensionPoint qualifiedName="org.jetbrains.binaryRequestHandler" interface="org.jetbrains.ide.BinaryRequestHandler"/>
<extensionPoint qualifiedName="org.jetbrains.customPortServerManager" interface="org.jetbrains.ide.CustomPortServerManager"/>
<extensionPoint qualifiedName="org.jetbrains.jsonRpcDomain" beanClass="org.jetbrains.io.jsonRpc.JsonRpcDomainBean">
<with attribute="implementation" implements="java.lang.Object"/>
</extensionPoint>
</extensionPoints>
<extensions defaultExtensionNs="com.intellij">
@@ -33,5 +37,7 @@
<webServerPathHandler implementation="org.jetbrains.builtInWebServer.DefaultWebServerPathHandler" order="last"/>
<webServerFileHandler implementation="org.jetbrains.builtInWebServer.BuiltInWebServer$StaticFileHandler" order="last"/>
<webServerRootsProvider implementation="org.jetbrains.builtInWebServer.DefaultWebServerRootsProvider"/>
<binaryRequestHandler implementation="org.jetbrains.io.jsonRpc.socket.RpcBinaryRequestHandler"/>
</extensions>
</idea-plugin>
@@ -48,6 +48,7 @@ public abstract class EvaluateContextBase<VALUE_MANAGER extends ValueManager> im
@NotNull
@Override
public Promise<?> refreshOnDone(@NotNull Promise<?> promise) {
//noinspection unchecked
return promise.then(valueManager.getClearCachesTask());
}
}
@@ -11,6 +11,7 @@ import org.jetbrains.concurrency.PromiseManager;
public abstract class ScriptManagerBase<SCRIPT extends ScriptBase> implements ScriptManager {
@SuppressWarnings("unchecked")
private final PromiseManager<ScriptBase, String> scriptSourceLoader = new PromiseManager<ScriptBase, String>(ScriptBase.class) {
@NotNull
@Override
public Promise<String> load(@NotNull ScriptBase script) {
//noinspection unchecked
@@ -17,6 +17,7 @@ public abstract class VariablesHost<VALUE_MANAGER extends ValueManager> {
return host.valueManager.getCacheStamp() == host.cacheStamp;
}
@NotNull
@Override
public Promise<List<Variable>> load(@NotNull VariablesHost host) {
if (host.valueManager.isObsolete()) {
@@ -1,8 +1,7 @@
package org.jetbrains.debugger.values;
import com.intellij.openapi.util.ActionCallback;
import com.intellij.util.Function;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.concurrency.ConsumerRunnable;
import org.jetbrains.concurrency.Promise;
import org.jetbrains.debugger.Vm;
@@ -33,11 +32,12 @@ public abstract class ValueManager<VM extends Vm> {
}
@NotNull
public ConsumerRunnable getClearCachesTask() {
return new ConsumerRunnable() {
public Function getClearCachesTask() {
return new Function<Object, Void>() {
@Override
public void run() {
public Void fun(Object o) {
clearCaches();
return null;
}
};
}
@@ -54,14 +54,6 @@ public abstract class ValueManager<VM extends Vm> {
obsolete = true;
}
public final boolean rejectIfObsolete(@NotNull ActionCallback result) {
if (isObsolete()) {
result.reject("Obsolete context");
return true;
}
return false;
}
@NotNull
public static <T> Promise<T> reject() {
//noinspection unchecked
@@ -19,7 +19,6 @@ import com.intellij.xdebugger.frame.presentation.XStringValuePresentation;
import com.intellij.xdebugger.frame.presentation.XValuePresentation;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.concurrency.ConsumerRunnable;
import org.jetbrains.concurrency.ObsolescentAsyncFunction;
import org.jetbrains.concurrency.Promise;
import org.jetbrains.debugger.values.*;
@@ -489,13 +488,16 @@ public final class VariableView extends XNamedValue implements VariableContext {
public void setValue(@NotNull String expression, @NotNull final XModificationCallback callback) {
ValueModifier valueModifier = variable.getValueModifier();
assert valueModifier != null;
valueModifier.setValue(variable, expression, getEvaluateContext()).done(new ConsumerRunnable() {
@Override
public void run() {
value = null;
callback.valueModified();
}
}).rejected(createErrorMessageConsumer(callback));
//noinspection unchecked
valueModifier.setValue(variable, expression, getEvaluateContext())
.done(new Consumer() {
@Override
public void consume(Object o) {
value = null;
callback.valueModified();
}
})
.rejected(createErrorMessageConsumer(callback));
}
};
}
@@ -621,14 +623,16 @@ public final class VariableView extends XNamedValue implements VariableContext {
}
final AtomicBoolean evaluated = new AtomicBoolean();
((StringValue)value).getFullString().done(new ConsumerRunnable() {
@Override
public void run() {
if (!callback.isObsolete() && evaluated.compareAndSet(false, true)) {
callback.evaluated(value.getValueString());
((StringValue)value).getFullString()
.done(new Consumer<String>() {
@Override
public void consume(String s) {
if (!callback.isObsolete() && evaluated.compareAndSet(false, true)) {
callback.evaluated(value.getValueString());
}
}
}
}).rejected(createErrorMessageConsumer(callback));
})
.rejected(createErrorMessageConsumer(callback));
}
}
@@ -48,7 +48,7 @@ public abstract class VmConnection<T extends Vm> implements Disposable, BrowserC
}
@NotNull
public AsyncPromise<Void> opened() {
public Promise<Void> opened() {
return opened;
}
@@ -19,9 +19,9 @@ import com.intellij.openapi.util.text.StringUtil;
import org.jetbrains.annotations.NotNull;
public class CharArrayCharSequence implements CharSequenceBackedByArray {
private final char[] myChars;
private final int myStart;
private final int myEnd;
protected final char[] myChars;
protected final int myStart;
protected final int myEnd;
public CharArrayCharSequence(@NotNull char... chars) {
this(chars, 0, chars.length);