-
Notifications
You must be signed in to change notification settings - Fork 1.6k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
* add new method isWritable to check netty channel state * add future based callback
- Loading branch information
sergeygrigorev
authored and
sergeygrigorev
committed
Oct 18, 2017
1 parent
567ad95
commit f844299
Showing
12 changed files
with
337 additions
and
32 deletions.
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
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
26 changes: 26 additions & 0 deletions
26
src/main/java/com/corundumstudio/socketio/NetworkCallback.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,26 @@ | ||
/** | ||
* Copyright 2012 Nikita Koksharov | ||
* | ||
* 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 com.corundumstudio.socketio; | ||
|
||
import java.util.concurrent.Future; | ||
|
||
/** | ||
* Callback for the network operations. | ||
*/ | ||
public interface NetworkCallback<V> extends Future<V> { | ||
/* todo: replace with <? extends Future<? super V>> */ | ||
void addCallback(NetworkCallbackListener<Future<? super V>> callback); | ||
} |
26 changes: 26 additions & 0 deletions
26
src/main/java/com/corundumstudio/socketio/NetworkCallbackListener.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,26 @@ | ||
/** | ||
* Copyright 2012 Nikita Koksharov | ||
* | ||
* 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 com.corundumstudio.socketio; | ||
|
||
import java.util.concurrent.Future; | ||
import java.util.EventListener; | ||
|
||
/** | ||
* Listener for {@link NetworkCallback}. | ||
*/ | ||
public interface NetworkCallbackListener<F extends Future<?>> extends EventListener { | ||
void operationComplete(F var1) throws Exception; | ||
} |
39 changes: 39 additions & 0 deletions
39
src/main/java/com/corundumstudio/socketio/NetworkCallbacks.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,39 @@ | ||
/** | ||
* Copyright 2012 Nikita Koksharov | ||
* | ||
* 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 com.corundumstudio.socketio; | ||
|
||
import com.corundumstudio.socketio.concurrent.FailedFuture; | ||
import com.corundumstudio.socketio.concurrent.SucceededFuture; | ||
import io.netty.util.concurrent.ImmediateEventExecutor; | ||
|
||
/** | ||
* Object instances of {@link NetworkCallback} | ||
*/ | ||
public class NetworkCallbacks { | ||
private static Throwable CHANNEL_CLOSED_ERROR = new IllegalStateException("channel is closed"); | ||
|
||
public static <T> NetworkCallback<T> channelClosed() { | ||
return new FailedFuture<T>(ImmediateEventExecutor.INSTANCE, CHANNEL_CLOSED_ERROR); | ||
} | ||
|
||
public static <T> NetworkCallback<T> success(T value) { | ||
return new SucceededFuture<T>(ImmediateEventExecutor.INSTANCE, value); | ||
} | ||
|
||
public static <E extends Throwable> NetworkCallback<E> failure(E throwable) { | ||
return new FailedFuture<E>(ImmediateEventExecutor.INSTANCE, throwable); | ||
} | ||
} |
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
43 changes: 43 additions & 0 deletions
43
src/main/java/com/corundumstudio/socketio/concurrent/CompleteFuture.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,43 @@ | ||
/** | ||
* Copyright 2012 Nikita Koksharov | ||
* | ||
* 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 com.corundumstudio.socketio.concurrent; | ||
|
||
import com.corundumstudio.socketio.NetworkCallback; | ||
import com.corundumstudio.socketio.NetworkCallbackListener; | ||
import io.netty.util.concurrent.EventExecutor; | ||
import io.netty.util.concurrent.GenericFutureListener; | ||
|
||
import java.util.concurrent.Future; | ||
|
||
/** | ||
* Simple promise which already has a result. | ||
*/ | ||
public abstract class CompleteFuture<V> extends io.netty.util.concurrent.CompleteFuture <V> implements NetworkCallback<V> { | ||
|
||
protected CompleteFuture(EventExecutor executor) { | ||
super(executor); | ||
} | ||
|
||
@Override | ||
public void addCallback(final NetworkCallbackListener<Future<? super V>> callback) { | ||
super.addListener(new GenericFutureListener<io.netty.util.concurrent.Future<? super V>>() { | ||
@Override | ||
public void operationComplete(io.netty.util.concurrent.Future<? super V> future) throws Exception { | ||
callback.operationComplete(future); | ||
} | ||
}); | ||
} | ||
} |
44 changes: 44 additions & 0 deletions
44
src/main/java/com/corundumstudio/socketio/concurrent/DefaultNetworkPromise.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,44 @@ | ||
/** | ||
* Copyright 2012 Nikita Koksharov | ||
* | ||
* 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 com.corundumstudio.socketio.concurrent; | ||
|
||
import com.corundumstudio.socketio.NetworkCallback; | ||
import com.corundumstudio.socketio.NetworkCallbackListener; | ||
import io.netty.util.concurrent.DefaultPromise; | ||
import io.netty.util.concurrent.EventExecutor; | ||
import io.netty.util.concurrent.GenericFutureListener; | ||
|
||
import java.util.concurrent.Future; | ||
|
||
/** | ||
* Simple network callback. | ||
*/ | ||
public class DefaultNetworkPromise<V> extends DefaultPromise<V> implements NetworkCallback<V> { | ||
|
||
public DefaultNetworkPromise(EventExecutor executor) { | ||
super(executor); | ||
} | ||
|
||
@Override | ||
public void addCallback(final NetworkCallbackListener<Future<? super V>> callback) { | ||
super.addListener(new GenericFutureListener<io.netty.util.concurrent.Future<? super V>>() { | ||
@Override | ||
public void operationComplete(io.netty.util.concurrent.Future<? super V> future) throws Exception { | ||
callback.operationComplete(future); | ||
} | ||
}); | ||
} | ||
} |
55 changes: 55 additions & 0 deletions
55
src/main/java/com/corundumstudio/socketio/concurrent/FailedFuture.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,55 @@ | ||
/** | ||
* Copyright 2012 Nikita Koksharov | ||
* | ||
* 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 com.corundumstudio.socketio.concurrent; | ||
|
||
import io.netty.util.concurrent.EventExecutor; | ||
import io.netty.util.concurrent.Future; | ||
import io.netty.util.internal.PlatformDependent; | ||
|
||
public final class FailedFuture<V> extends CompleteFuture<V> { | ||
private final Throwable cause; | ||
|
||
public FailedFuture(EventExecutor executor, Throwable cause) { | ||
super(executor); | ||
if (cause == null) { | ||
throw new NullPointerException("cause"); | ||
} else { | ||
this.cause = cause; | ||
} | ||
} | ||
|
||
public Throwable cause() { | ||
return this.cause; | ||
} | ||
|
||
public boolean isSuccess() { | ||
return false; | ||
} | ||
|
||
public Future<V> sync() { | ||
PlatformDependent.throwException(this.cause); | ||
return this; | ||
} | ||
|
||
public Future<V> syncUninterruptibly() { | ||
PlatformDependent.throwException(this.cause); | ||
return this; | ||
} | ||
|
||
public V getNow() { | ||
return null; | ||
} | ||
} |
39 changes: 39 additions & 0 deletions
39
src/main/java/com/corundumstudio/socketio/concurrent/SucceededFuture.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,39 @@ | ||
/** | ||
* Copyright 2012 Nikita Koksharov | ||
* | ||
* 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 com.corundumstudio.socketio.concurrent; | ||
|
||
import io.netty.util.concurrent.EventExecutor; | ||
|
||
public final class SucceededFuture<V> extends CompleteFuture<V> { | ||
private final V result; | ||
|
||
public SucceededFuture(EventExecutor executor, V result) { | ||
super(executor); | ||
this.result = result; | ||
} | ||
|
||
public Throwable cause() { | ||
return null; | ||
} | ||
|
||
public boolean isSuccess() { | ||
return true; | ||
} | ||
|
||
public V getNow() { | ||
return this.result; | ||
} | ||
} |
Oops, something went wrong.