use ByteBuf directly — avoid intermediate StringBuilder

We can create efficient OutputStreamWriter backed by ByteBuf, but it will be implemented later
This commit is contained in:
Vladimir Krivosheev
2015-02-16 09:18:44 +01:00
parent f7308adb67
commit 950adb3dff
10 changed files with 244 additions and 151 deletions
@@ -1,7 +1,6 @@
package org.jetbrains.io.jsonRpc;
import com.google.gson.*;
import com.google.gson.internal.Streams;
import com.google.gson.reflect.TypeToken;
import com.google.gson.stream.JsonToken;
import com.google.gson.stream.JsonWriter;
@@ -9,6 +8,7 @@ import com.intellij.openapi.diagnostic.Logger;
import com.intellij.openapi.util.NotNullLazyValue;
import com.intellij.openapi.util.Pair;
import com.intellij.openapi.util.Ref;
import com.intellij.openapi.vfs.CharsetToolkit;
import com.intellij.util.ArrayUtil;
import com.intellij.util.ArrayUtilRt;
import com.intellij.util.Consumer;
@@ -17,10 +17,7 @@ import com.intellij.util.text.CharSequenceBackedByArray;
import gnu.trove.THashMap;
import gnu.trove.TIntArrayList;
import gnu.trove.TIntProcedure;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.ByteBufAllocator;
import io.netty.buffer.ByteBufUtil;
import io.netty.util.CharsetUtil;
import io.netty.buffer.*;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.concurrency.AsyncPromise;
@@ -29,9 +26,10 @@ import org.jetbrains.io.JsonReaderEx;
import org.jetbrains.io.JsonUtil;
import java.io.IOException;
import java.io.OutputStreamWriter;
import java.lang.reflect.Method;
import java.lang.reflect.Type;
import java.nio.CharBuffer;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
@@ -171,12 +169,12 @@ public class JsonRpcServer implements MessageServer {
}
}
public void sendResponse(int messageId, @NotNull Client client, @Nullable CharSequence rawMessage) {
public void sendResponse(int messageId, @NotNull Client client, @Nullable ByteBuf rawMessage) {
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(client.getByteBufAllocator(), messageId, "e", null, rawMessage, ArrayUtil.EMPTY_OBJECT_ARRAY));
client.send(encodeMessage(client.getByteBufAllocator(), messageId, "e", null, null, new Object[]{rawMessage}));
}
@SuppressWarnings("unused")
@@ -194,7 +192,7 @@ public class JsonRpcServer implements MessageServer {
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) {
public boolean sendWithRawPart(@NotNull Client client, @NotNull String domain, @NotNull String command, @Nullable ByteBuf rawMessage, Object... params) {
client.send(encodeMessage(client.getByteBufAllocator(), -1, domain, command, rawMessage, params));
return true;
}
@@ -217,58 +215,100 @@ public class JsonRpcServer implements MessageServer {
int messageId,
@Nullable String domain,
@Nullable String command,
@Nullable CharSequence rawData,
@Nullable ByteBuf rawData,
@Nullable Object[] params) {
StringBuilder sb = new StringBuilder();
ByteBuf buffer = byteBufAllocator.ioBuffer();
boolean success = false;
try {
doEncodeMessage(sb, messageId, domain, command, params == null ? ArrayUtil.EMPTY_OBJECT_ARRAY : params, rawData);
Object[] notNullParams = params == null ? ArrayUtil.EMPTY_OBJECT_ARRAY : params;
buffer = doEncodeMessage(byteBufAllocator, buffer, messageId, domain, command, notNullParams, rawData);
if (LOG.isDebugEnabled()) {
LOG.debug("OUT " + sb.toString());
LOG.debug("OUT " + domain + '.' + command + (notNullParams.length == 0 ? "" : " " + Arrays.toString(params)) + (rawData == null ? "" : " " + rawData.toString(CharsetToolkit.UTF8_CHARSET)));
}
return ByteBufUtil.encodeString(byteBufAllocator, CharBuffer.wrap(sb), CharsetUtil.UTF_8);
success = true;
return buffer;
}
catch (IOException e) {
throw new RuntimeException(e);
}
finally {
if (!success) {
buffer.release();
}
}
}
private void doEncodeMessage(@NotNull StringBuilder sb, int id, @Nullable String domain, @Nullable String command, @NotNull Object[] params, @Nullable CharSequence rawData) throws IOException {
sb.append('[');
@NotNull
private ByteBuf doEncodeMessage(@NotNull ByteBufAllocator byteBufAllocator,
@NotNull ByteBuf buffer,
int id,
@Nullable String domain,
@Nullable String command,
@NotNull Object[] params,
@Nullable ByteBuf rawData) throws IOException {
buffer.writeByte('[');
ByteBuf effectiveBuffer = buffer;
boolean hasPrev = false;
StringBuilder sb = null;
if (id != -1) {
sb.append(id);
sb = new StringBuilder();
ByteBufUtil.writeAscii(buffer, sb.append(id));
sb.setLength(0);
hasPrev = true;
}
if (domain != null) {
if (hasPrev) {
sb.append(',');
buffer.writeByte(',');
}
sb.append('"').append(domain).append("\",\"");
buffer.writeByte('"');
ByteBufUtil.writeAscii(buffer, domain);
buffer.writeByte('"').writeByte(',').writeByte('"');
if (command == null) {
if (rawData != null) {
sb.append(rawData);
effectiveBuffer = byteBufAllocator.compositeBuffer().addComponent(buffer).addComponent(rawData);
buffer = byteBufAllocator.ioBuffer();
}
sb.append('"');
return;
buffer.writeByte('"');
return addBuffer(effectiveBuffer, buffer);
}
else {
sb.append(command).append('"');
ByteBufUtil.writeAscii(buffer, command);
buffer.writeByte('"');
}
}
encodeParameters(sb, params, rawData);
sb.append(']');
encodeParameters(buffer, params, rawData, sb);
if (rawData != null) {
if (params.length > 0) {
buffer.writeByte(',');
}
effectiveBuffer = byteBufAllocator.compositeBuffer().addComponent(buffer).addComponent(rawData);
buffer = byteBufAllocator.ioBuffer();
}
buffer.writeByte(']');
buffer.writeByte(']');
return addBuffer(effectiveBuffer, buffer);
}
private void encodeParameters(@NotNull StringBuilder sb, @NotNull Object[] params, @Nullable CharSequence rawData) throws IOException {
@NotNull
// addComponent always add sliced component, so, we must add last buffer only after all writes
private static ByteBuf addBuffer(@NotNull ByteBuf buffer, @NotNull ByteBuf lastComponent) {
if (buffer != lastComponent) {
((CompositeByteBuf)buffer).addComponent(lastComponent);
buffer.writerIndex(buffer.capacity());
}
return buffer;
}
private void encodeParameters(@NotNull ByteBuf buffer, @NotNull Object[] params, @Nullable ByteBuf rawData, @Nullable StringBuilder sb) throws IOException {
JsonWriter writer = null;
sb.append(',').append('[');
buffer.writeByte(',').writeByte('[');
boolean hasPrev = false;
for (Object param : params) {
if (hasPrev) {
sb.append(',');
buffer.writeByte(',');
}
else {
hasPrev = true;
@@ -276,34 +316,53 @@ public class JsonRpcServer implements MessageServer {
// gson - SOE if param has type class com.intellij.openapi.editor.impl.DocumentImpl$MyCharArray, so, use hack
if (param instanceof CharSequence) {
JsonUtil.escape(((CharSequence)param), sb);
JsonUtil.escape(((CharSequence)param), buffer);
}
else if (param == null) {
sb.append("null");
ByteBufUtil.writeAscii(buffer, "null");
}
else if (param instanceof Number || param instanceof Boolean) {
sb.append(param.toString());
else if (param instanceof Boolean) {
ByteBufUtil.writeAscii(buffer, param.toString());
}
else if (param instanceof Number) {
if (sb == null) {
sb = new StringBuilder();
}
if (param instanceof Integer) {
sb.append(((Integer)param).intValue());
}
else if (param instanceof Long) {
sb.append(((Long)param).longValue());
}
else if (param instanceof Float) {
sb.append(((Float)param).floatValue());
}
else if (param instanceof Double) {
sb.append(((Double)param).doubleValue());
}
else {
sb.append(param.toString());
}
ByteBufUtil.writeAscii(buffer, sb);
sb.setLength(0);
}
else if (param instanceof Consumer) {
if (sb == null) {
sb = new StringBuilder();
}
//noinspection unchecked
((Consumer<StringBuilder>)param).consume(sb);
ByteBufUtil.writeUtf8(buffer, sb);
sb.setLength(0);
}
else {
if (writer == null) {
writer = new JsonWriter(Streams.writerForAppendable(sb));
writer = new JsonWriter(new OutputStreamWriter(new ByteBufOutputStream(buffer)));
}
//noinspection unchecked
((TypeAdapter<Object>)gson.getAdapter(param.getClass())).write(writer, param);
}
}
if (rawData != null) {
if (hasPrev) {
sb.append(',');
}
sb.append(rawData);
}
sb.append(']');
}
private static class IntArrayListTypeAdapter<T> extends TypeAdapter<T> {
@@ -0,0 +1,68 @@
/*
* 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 io.netty.buffer;
import io.netty.util.CharsetUtil;
public class ByteBufUtilEx {
// todo pull request
public static int writeUtf8(ByteBuf buf, CharSequence seq, int start, int end) {
if (buf == null) {
throw new NullPointerException("buf");
}
if (seq == null) {
throw new NullPointerException("seq");
}
// UTF-8 uses max. 3 bytes per char, so calculate the worst case.
final int len = end - start;
final int maxSize = len * 3;
buf.ensureWritable(maxSize);
if (buf instanceof AbstractByteBuf) {
// Fast-Path
AbstractByteBuf buffer = (AbstractByteBuf)buf;
int oldWriterIndex = buffer.writerIndex;
int writerIndex = oldWriterIndex;
// We can use the _set methods as these not need to do any index checks and reference checks.
// This is possible as we called ensureWritable(...) before.
for (int i = start; i < end; i++) {
char c = seq.charAt(i);
if (c < 0x80) {
buffer._setByte(writerIndex++, (byte)c);
}
else if (c < 0x800) {
buffer._setByte(writerIndex++, (byte)(0xc0 | (c >> 6)));
buffer._setByte(writerIndex++, (byte)(0x80 | (c & 0x3f)));
}
else {
buffer._setByte(writerIndex++, (byte)(0xe0 | (c >> 12)));
buffer._setByte(writerIndex++, (byte)(0x80 | ((c >> 6) & 0x3f)));
buffer._setByte(writerIndex++, (byte)(0x80 | (c & 0x3f)));
}
}
// update the writerIndex without any extra checks for performance reasons
buffer.writerIndex = writerIndex;
return writerIndex - oldWriterIndex;
}
else {
// Maybe we could also check if we can unwrap() to access the wrapped buffer which
// may be an AbstractByteBuf. But this may be overkill so let us keep it simple for now.
byte[] bytes = seq.toString().getBytes(CharsetUtil.UTF_8);
buf.writeBytes(bytes);
return bytes.length;
}
}
}
@@ -3,6 +3,9 @@ package org.jetbrains.io;
import com.intellij.util.ArrayUtil;
import com.intellij.util.SmartList;
import gnu.trove.THashMap;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.ByteBufUtil;
import io.netty.buffer.ByteBufUtilEx;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
@@ -62,6 +65,41 @@ public class JsonUtil {
sb.append('"');
}
public static void escape(@NotNull CharSequence value, @NotNull ByteBuf buffer) {
int length = value.length();
buffer.ensureWritable(length * 2);
buffer.writeByte('"');
int last = 0;
for (int i = 0; i < length; i++) {
char c = value.charAt(i);
String replacement;
if (c < 128) {
replacement = REPLACEMENT_CHARS[c];
if (replacement == null) {
continue;
}
}
else if (c == '\u2028') {
replacement = "\\u2028";
}
else if (c == '\u2029') {
replacement = "\\u2029";
}
else {
continue;
}
if (last < i) {
ByteBufUtilEx.writeUtf8(buffer, value, last, i);
}
ByteBufUtil.writeAscii(buffer, replacement);
last = i + 1;
}
if (last < length) {
ByteBufUtilEx.writeUtf8(buffer, value, last, length);
}
buffer.writeByte('"');
}
@NotNull
public static <T> List<T> nextList(@NotNull JsonReaderEx reader) {
reader.beginArray();
@@ -29,6 +29,8 @@ 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;
@@ -217,6 +219,6 @@ public final class NettyUtil {
if (pipeline.get(ChunkedWriteHandler.class) == null) {
pipeline.addLast("chunkedWriteHandler", new ChunkedWriteHandler());
}
//pipeline.addLast("corsHandler", new CorsHandler(CorsConfig.withAnyOrigin().allowCredentials().allowNullOrigin().allowedRequestMethods().build()));
pipeline.addLast("corsHandler", new CorsHandler(CorsConfig.withAnyOrigin().allowCredentials().allowNullOrigin().allowedRequestMethods().build()));
}
}
@@ -1,6 +1,7 @@
package org.jetbrains.debugger;
import com.intellij.util.Consumer;
import io.netty.buffer.ByteBuf;
import io.netty.channel.Channel;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelFutureListener;
@@ -23,7 +24,7 @@ public class StandaloneVmHelper extends MessageWriter implements Vm.AttachStateM
}
@Override
public boolean write(@NotNull CharSequence content) {
public boolean write(@NotNull ByteBuf content) {
return write(((Object)content));
}
@@ -1,6 +1,8 @@
package org.jetbrains.rpc;
import com.intellij.openapi.vfs.CharsetToolkit;
import com.intellij.util.BooleanFunction;
import io.netty.buffer.ByteBuf;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.jsonProtocol.Request;
@@ -9,9 +11,9 @@ import static org.jetbrains.rpc.CommandProcessor.LOG;
public abstract class MessageWriter implements BooleanFunction<Request> {
@Override
public boolean fun(@NotNull Request message) {
CharSequence content = message.toJson();
ByteBuf content = message.getBuffer();
if (isDebugLoggingEnabled()) {
LOG.debug("OUT: " + content.toString());
LOG.debug("OUT: " + content.toString(CharsetToolkit.UTF8_CHARSET));
}
return write(content);
}
@@ -20,5 +22,5 @@ public abstract class MessageWriter implements BooleanFunction<Request> {
return LOG.isDebugEnabled();
}
protected abstract boolean write(@NotNull CharSequence content);
protected abstract boolean write(@NotNull ByteBuf content);
}
@@ -10,5 +10,6 @@
<orderEntry type="library" exported="" name="gson" level="project" />
<orderEntry type="module" module-name="util" />
<orderEntry type="module" module-name="platform-impl" />
<orderEntry type="library" name="Netty" level="project" />
</component>
</module>
@@ -1,22 +1,28 @@
package org.jetbrains.jsonProtocol;
import com.google.gson.stream.JsonWriter;
import com.intellij.openapi.vfs.CharsetToolkit;
import gnu.trove.TIntArrayList;
import gnu.trove.TIntHashSet;
import gnu.trove.TIntProcedure;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.ByteBufAllocator;
import io.netty.buffer.ByteBufOutputStream;
import io.netty.buffer.ByteBufUtil;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.io.JsonUtil;
import java.io.IOException;
import java.io.OutputStreamWriter;
import java.util.Collection;
import java.util.List;
import java.util.Map;
public abstract class OutMessage {
@SuppressWarnings("IOResourceOpenedButNotSafelyClosed")
private final StringWriter stringWriter = new StringWriter();
public final JsonWriter writer = new JsonWriter(stringWriter);
// todo don't want to risk - are we really can release it properly? heap buffer for now
private final ByteBuf buffer = ByteBufAllocator.DEFAULT.heapBuffer();
public final JsonWriter writer = new JsonWriter(new OutputStreamWriter(new ByteBufOutputStream(buffer)));
private boolean finalized;
@@ -177,13 +183,12 @@ public abstract class OutMessage {
boolean isNotFirst = false;
for (OutMessage item : value) {
if (isNotFirst) {
stringWriter.append(',').append(' ');
buffer.writeByte(',').writeByte(' ');
}
else {
isNotFirst = true;
}
StringBuilder buffer = item.stringWriter.getBuffer();
if (!item.finalized) {
item.finalized = true;
try {
@@ -191,7 +196,7 @@ public abstract class OutMessage {
}
catch (IllegalStateException e) {
if ("Nesting problem.".equals(e.getMessage())) {
throw new RuntimeException(item.stringWriter.getBuffer() + "\nparent:\n" + stringWriter.getBuffer(), e);
throw new RuntimeException(item.buffer.toString(CharsetToolkit.UTF8_CHARSET) + "\nparent:\n" + buffer.toString(CharsetToolkit.UTF8_CHARSET), e);
}
else {
throw e;
@@ -199,7 +204,7 @@ public abstract class OutMessage {
}
}
stringWriter.append(buffer);
buffer.writeBytes(item.buffer);
}
writer.endArray();
}
@@ -218,26 +223,26 @@ public abstract class OutMessage {
}
}
public static void prepareWriteRaw(OutMessage message, String name) throws IOException {
public static void prepareWriteRaw(@NotNull OutMessage message, @NotNull String name) throws IOException {
message.writer.name(name).nullValue();
StringBuilder myBuffer = message.stringWriter.getBuffer();
myBuffer.delete(myBuffer.length() - "null".length(), myBuffer.length());
message.writer.flush();
ByteBuf itemBuffer = message.buffer;
itemBuffer.writerIndex(itemBuffer.writerIndex() - "null".length());
}
public static void doWriteRaw(OutMessage message, String rawValue) {
message.stringWriter.append(rawValue);
public static void doWriteRaw(@NotNull OutMessage message, @NotNull String rawValue) {
ByteBufUtil.writeUtf8(message.buffer, rawValue);
}
protected final void writeMessage(String name, OutMessage value) {
protected final void writeMessage(@NotNull String name, @NotNull OutMessage value) {
try {
beginArguments();
prepareWriteRaw(this, name);
StringBuilder buffer = value.stringWriter.getBuffer();
if (!value.finalized) {
value.close();
}
stringWriter.append(buffer);
buffer.writeBytes(value.buffer);
}
catch (IOException e) {
throw new RuntimeException(e);
@@ -287,11 +292,11 @@ public abstract class OutMessage {
}
}
protected final void writeString(String name, CharSequence value) {
protected final void writeString(@NotNull String name, CharSequence value) {
if (value != null) {
try {
prepareWriteRaw(this, name);
JsonUtil.escape(value, stringWriter.getBuffer());
JsonUtil.escape(value, buffer);
}
catch (IOException e) {
throw new RuntimeException(e);
@@ -311,7 +316,7 @@ public abstract class OutMessage {
@NotNull
@SuppressWarnings("UnusedDeclaration")
public final CharSequence toJson() {
return stringWriter.getBuffer();
public final ByteBuf getBuffer() {
return buffer;
}
}
@@ -1,10 +1,11 @@
package org.jetbrains.jsonProtocol;
import io.netty.buffer.ByteBuf;
import org.jetbrains.annotations.NotNull;
public interface Request<RESULT> {
@NotNull
CharSequence toJson();
ByteBuf getBuffer();
String getMethodName();
@@ -1,84 +0,0 @@
package org.jetbrains.jsonProtocol;
import org.jetbrains.annotations.NotNull;
import java.io.IOException;
import java.io.Writer;
public class StringWriter extends Writer {
private final StringBuilder builder;
public StringWriter() {
builder = new StringBuilder();
lock = builder;
}
public StringWriter(int initialSize) {
if (initialSize < 0) {
throw new IllegalArgumentException("Negative buffer size");
}
builder = new StringBuilder(initialSize);
lock = builder;
}
@Override
public void write(int c) {
builder.append((char)c);
}
@Override
public void write(@NotNull char[] chars, int off, int len) {
builder.append(chars, off, len);
}
@Override
public void write(@NotNull String string) {
builder.append(string);
}
@Override
public void write(@NotNull String str, int off, int length) {
builder.append(str, off, off + length);
}
@NotNull
@Override
public StringWriter append(CharSequence charSequence) {
builder.append(charSequence);
return this;
}
@NotNull
@Override
public StringWriter append(CharSequence charSequence, int start, int end) {
builder.append(charSequence, start, end);
return this;
}
@NotNull
@Override
public StringWriter append(char c) {
write(c);
return this;
}
public String toString() {
return builder.toString();
}
public StringBuilder getBuffer() {
return builder;
}
/**
* Flush the stream.
*/
@Override
public void flush() {
}
@Override
public void close() throws IOException {
}
}