Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Nikita
committed
Jul 8, 2015
1 parent
6c182ed
commit e5da696
Showing
16 changed files
with
371 additions
and
89 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
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
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
44 changes: 44 additions & 0 deletions
44
src/main/java/org/redisson/client/RedisPubSubConnection.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 @@ | ||
package org.redisson.client; | ||
|
||
import org.redisson.client.handler.RedisData; | ||
import org.redisson.client.protocol.Codec; | ||
import org.redisson.client.protocol.PubSubMessage; | ||
import org.redisson.client.protocol.PubSubMessageDecoder; | ||
import org.redisson.client.protocol.RedisCommand; | ||
import org.redisson.client.protocol.RedisCommands; | ||
import org.redisson.client.protocol.StringCodec; | ||
|
||
import io.netty.channel.Channel; | ||
import io.netty.channel.ChannelFuture; | ||
import io.netty.util.concurrent.Future; | ||
import io.netty.util.concurrent.Promise; | ||
|
||
public class RedisPubSubConnection { | ||
|
||
final Channel channel; | ||
final RedisClient redisClient; | ||
|
||
public RedisPubSubConnection(RedisClient redisClient, Channel channel) { | ||
this.redisClient = redisClient; | ||
this.channel = channel; | ||
} | ||
|
||
public Future<PubSubMessage> subscribe(String ... channel) { | ||
return async(new PubSubMessageDecoder(), RedisCommands.SUBSCRIBE, channel); | ||
} | ||
|
||
public Future<Long> publish(String channel, String msg) { | ||
return async(new StringCodec(), RedisCommands.PUBLISH, channel, msg); | ||
} | ||
|
||
public <T, R> Future<R> async(Codec encoder, RedisCommand<T> command, Object ... params) { | ||
Promise<R> promise = redisClient.getBootstrap().group().next().<R>newPromise(); | ||
channel.writeAndFlush(new RedisData<T, R>(promise, encoder, command, params)); | ||
return promise; | ||
} | ||
|
||
public ChannelFuture closeAsync() { | ||
return channel.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
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
29 changes: 29 additions & 0 deletions
29
src/main/java/org/redisson/client/protocol/PubSubMessage.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,29 @@ | ||
package org.redisson.client.protocol; | ||
|
||
public class PubSubMessage { | ||
|
||
public enum Type {SUBSCRIBE, MESSAGE} | ||
|
||
private Type type; | ||
private String channel; | ||
|
||
public PubSubMessage(Type type, String channel) { | ||
super(); | ||
this.type = type; | ||
this.channel = channel; | ||
} | ||
|
||
public String getChannel() { | ||
return channel; | ||
} | ||
|
||
public Type getType() { | ||
return type; | ||
} | ||
|
||
@Override | ||
public String toString() { | ||
return "PubSubReplay [type=" + type + ", channel=" + channel + "]"; | ||
} | ||
|
||
} |
Oops, something went wrong.