-
Notifications
You must be signed in to change notification settings - Fork 24.6k
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Add nio http server transport #29587
Merged
Merged
Changes from 74 commits
Commits
Show all changes
81 commits
Select commit
Hold shift + click to select a range
17c904a
Add server context
Tim-Brooks 6c8cd44
WIP
Tim-Brooks d9d995b
Do not extend autocloseable
Tim-Brooks 88bca9a
But keep close method
Tim-Brooks c4db506
Remove multiple context getters
Tim-Brooks 4c3bf37
WIP
Tim-Brooks dfd1901
Merge remote-tracking branch 'upstream/master' into layer_http
Tim-Brooks c4e1c50
Merge branch 'master' into layer_http
Tim-Brooks b5a374a
Do not depend on netty module
Tim-Brooks 408c87c
Pull over tests
Tim-Brooks 0f74aa0
Work on bytes producer
Tim-Brooks 44be919
Merge remote-tracking branch 'upstream/master' into layer_http
Tim-Brooks 913028e
Comments
Tim-Brooks df53494
Merge remote-tracking branch 'upstream/master' into layer_http
Tim-Brooks 9072a45
WIP
Tim-Brooks 5c89990
Merge remote-tracking branch 'upstream/master' into layer_http
Tim-Brooks 8bccda3
WIP
Tim-Brooks 26bdadf
Merge remote-tracking branch 'upstream/master' into layer_http
Tim-Brooks 99ad1fd
Work on refactoring
Tim-Brooks 3540e5c
Continue op refactor
Tim-Brooks bcc776a
WIP
Tim-Brooks 1e86182
Work on fixing tests
Tim-Brooks c6d01b7
Work on tests
Tim-Brooks 234e965
get simple tests passing
Tim-Brooks 4e9ab9d
WIP
Tim-Brooks 6b16068
Close write producer
Tim-Brooks c3c9c4f
Move bytes writer
Tim-Brooks fe4c961
Move producer
Tim-Brooks 8f5622d
Remove imports
Tim-Brooks e7fc228
Merge remote-tracking branch 'upstream/master' into layer_http
Tim-Brooks 0bd3810
Remove http stuff
Tim-Brooks 6bd2d91
Extract stuff to super
Tim-Brooks c761bce
Fix tests
Tim-Brooks e196257
Fix warning
Tim-Brooks 35daf28
Merge remote-tracking branch 'upstream/master' into remove_http
Tim-Brooks ac7ba90
Changes based on review
Tim-Brooks b1047b0
Work on netty adaptor
Tim-Brooks 10c75aa
Start http adapting work
Tim-Brooks 21dbd7e
Work on implementing http
Tim-Brooks 4fe6870
Work on server transport
Tim-Brooks 70f50c8
Work on config
Tim-Brooks 3e56692
WIP
Tim-Brooks a94d6f6
WIP
Tim-Brooks e0aa6bf
Work on setting up transport
Tim-Brooks 56b0cf5
Make basic http tests pass
Tim-Brooks 1c2e5d6
Wip
Tim-Brooks 86438d1
Fix issue
Tim-Brooks 7efb126
Fix checkstyle
Tim-Brooks a96dc7f
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks 8dfe123
Fix checkstyle
Tim-Brooks 9c77802
Fix forbidden
Tim-Brooks 30301b0
headers
Tim-Brooks dfdd446
Update setting
Tim-Brooks b3dc74c
Fix third party audit
Tim-Brooks eeeb8fd
Fix tests
Tim-Brooks e4380f2
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks 6a73065
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks c18f87f
Update
Tim-Brooks 4960e16
Fix issue
Tim-Brooks bb486b3
Fix compile issue
Tim-Brooks 0c7d64a
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks b29b80c
Cleanups
Tim-Brooks 94c748b
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks be40373
Move tests to abstract tests class
Tim-Brooks 996b557
Remove unnecessary methods
Tim-Brooks 683142c
WIP
Tim-Brooks ec56f62
Fix warning
Tim-Brooks 95debd9
Fix warning
Tim-Brooks bd81563
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks dd3f129
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks 6f64963
Close and wait
Tim-Brooks 29a942f
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks 2a40a6c
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks 839f071
Actualy close the channel
Tim-Brooks 570c6ed
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks f66801f
Work on read write tests
Tim-Brooks e2564c0
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks 52ef47e
Changes from review
Tim-Brooks c6d084e
A few more tests
Tim-Brooks 408b845
Remove unused imports
Tim-Brooks 3f01084
Merge remote-tracking branch 'upstream/master' into add_http_back
Tim-Brooks File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
49 changes: 49 additions & 0 deletions
49
libs/elasticsearch-nio/src/main/java/org/elasticsearch/nio/BytesWriteHandler.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,49 @@ | ||
/* | ||
* Licensed to Elasticsearch under one or more contributor | ||
* license agreements. See the NOTICE file distributed with | ||
* this work for additional information regarding copyright | ||
* ownership. Elasticsearch licenses this file to you 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.elasticsearch.nio; | ||
|
||
import java.nio.ByteBuffer; | ||
import java.util.Collections; | ||
import java.util.List; | ||
import java.util.function.BiConsumer; | ||
|
||
public abstract class BytesWriteHandler implements ReadWriteHandler { | ||
|
||
private static final List<FlushOperation> EMPTY_LIST = Collections.emptyList(); | ||
|
||
public WriteOperation createWriteOperation(SocketChannelContext context, Object message, BiConsumer<Void, Throwable> listener) { | ||
if (message instanceof ByteBuffer[]) { | ||
return new FlushReadyWrite(context, (ByteBuffer[]) message, listener); | ||
} else { | ||
throw new IllegalArgumentException("This channel only supports messages that are of type: " + ByteBuffer[].class); | ||
} | ||
} | ||
|
||
public List<FlushOperation> writeToBytes(WriteOperation writeOperation) { | ||
assert writeOperation instanceof FlushReadyWrite : "Write operation must be flush ready"; | ||
return Collections.singletonList((FlushReadyWrite) writeOperation); | ||
} | ||
|
||
public List<FlushOperation> pollFlushOperations() { | ||
return EMPTY_LIST; | ||
} | ||
|
||
public void close() {} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
45 changes: 45 additions & 0 deletions
45
libs/elasticsearch-nio/src/main/java/org/elasticsearch/nio/FlushReadyWrite.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,45 @@ | ||
/* | ||
* Licensed to Elasticsearch under one or more contributor | ||
* license agreements. See the NOTICE file distributed with | ||
* this work for additional information regarding copyright | ||
* ownership. Elasticsearch licenses this file to you 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.elasticsearch.nio; | ||
|
||
import java.nio.ByteBuffer; | ||
import java.util.function.BiConsumer; | ||
|
||
public class FlushReadyWrite extends FlushOperation implements WriteOperation { | ||
|
||
private final SocketChannelContext channelContext; | ||
private final ByteBuffer[] buffers; | ||
|
||
FlushReadyWrite(SocketChannelContext channelContext, ByteBuffer[] buffers, BiConsumer<Void, Throwable> listener) { | ||
super(buffers, listener); | ||
this.channelContext = channelContext; | ||
this.buffers = buffers; | ||
} | ||
|
||
@Override | ||
public SocketChannelContext getChannel() { | ||
return channelContext; | ||
} | ||
|
||
@Override | ||
public ByteBuffer[] getObject() { | ||
return buffers; | ||
} | ||
} |
71 changes: 71 additions & 0 deletions
71
libs/elasticsearch-nio/src/main/java/org/elasticsearch/nio/ReadWriteHandler.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,71 @@ | ||
/* | ||
* Licensed to Elasticsearch under one or more contributor | ||
* license agreements. See the NOTICE file distributed with | ||
* this work for additional information regarding copyright | ||
* ownership. Elasticsearch licenses this file to you 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.elasticsearch.nio; | ||
|
||
import java.io.IOException; | ||
import java.util.List; | ||
import java.util.function.BiConsumer; | ||
|
||
/** | ||
* Implements the application specific logic for handling inbound and outbound messages for a channel. | ||
*/ | ||
public interface ReadWriteHandler { | ||
|
||
/** | ||
* This method is called when a message is queued with a channel. It can be called from any thread. | ||
* This method should validate that the message is a valid type and return a write operation object | ||
* to be queued with the channel | ||
* | ||
* @param context the channel context | ||
* @param message the message | ||
* @param listener the listener to be called when the message is sent | ||
* @return the write operation to be queued | ||
*/ | ||
WriteOperation createWriteOperation(SocketChannelContext context, Object message, BiConsumer<Void, Throwable> listener); | ||
|
||
/** | ||
* This method is called on the event loop thread. It should serialize a write operation object to bytes | ||
* that can be flushed to the raw nio channel. | ||
* | ||
* @param writeOperation to be converted to bytes | ||
* @return the operations to flush the bytes to the channel | ||
*/ | ||
List<FlushOperation> writeToBytes(WriteOperation writeOperation); | ||
|
||
/** | ||
* Returns any flush operations that are ready to flush. This exists as a way to check if any flush | ||
* operations were produced during a read call. | ||
* | ||
* @return flush operations | ||
*/ | ||
List<FlushOperation> pollFlushOperations(); | ||
|
||
/** | ||
* This method handles bytes that have been read from the network. It should return the number of bytes | ||
* consumed so that they can be released. | ||
* | ||
* @param channelBuffer of bytes read from the network | ||
* @return the number of bytes consumed | ||
* @throws IOException if an exception occurs | ||
*/ | ||
int consumeReads(InboundChannelBuffer channelBuffer) throws IOException; | ||
|
||
void close() throws IOException; | ||
} |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
any chance we can make this typesafe?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
This is something I have thought about and I don't think it is very easy to do in this PR. The
ReadWriteHandler
would need to be parameterized. And then the channel context would need to be parameterized. And then the channel.I think down the line it might be more possible. With implementing an "http" channel type that is a subclass of a niochannel (similar to how we have a tcp channel type). But I think that would be a lot more code and better to do as follow-up work.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
++