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
黄志磊
committed
Dec 23, 2015
1 parent
45a31c7
commit 3ba747c
Showing
16 changed files
with
387 additions
and
79 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
21 changes: 21 additions & 0 deletions
21
mpush-connection/src/main/java/com/shinemo/mpush/connection/netty/NettySharedHolder.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,21 @@ | ||
package com.shinemo.mpush.connection.netty; | ||
|
||
import com.shinemo.mpush.core.thread.NamedThreadFactory; | ||
import com.shinemo.mpush.core.thread.ThreadNameSpace; | ||
|
||
import io.netty.buffer.ByteBufAllocator; | ||
import io.netty.buffer.UnpooledByteBufAllocator; | ||
import io.netty.util.HashedWheelTimer; | ||
import io.netty.util.Timer; | ||
|
||
public class NettySharedHolder { | ||
|
||
public static final Timer timer = new HashedWheelTimer(new NamedThreadFactory(ThreadNameSpace.NETTY_TIMER)); | ||
|
||
public static final ByteBufAllocator byteBufAllocator; | ||
|
||
static { | ||
byteBufAllocator = UnpooledByteBufAllocator.DEFAULT; | ||
} | ||
|
||
} |
2 changes: 1 addition & 1 deletion
2
...emo/mpush/api/protocol/PacketDecoder.java → ...nnection/netty/encoder/PacketDecoder.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
4 changes: 3 additions & 1 deletion
4
...emo/mpush/api/protocol/PacketEncoder.java → ...nnection/netty/encoder/PacketEncoder.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
2 changes: 1 addition & 1 deletion
2
...o/mpush/connection/ConnectionHandler.java → ...tion/netty/handler/ConnectionHandler.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
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
57 changes: 0 additions & 57 deletions
57
mpush-core/src/main/java/com/shinemo/mpush/core/ConnectionImpl.java
This file was deleted.
Oops, something went wrong.
54 changes: 45 additions & 9 deletions
54
mpush-core/src/main/java/com/shinemo/mpush/core/ConnectionManager.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 |
---|---|---|
@@ -1,27 +1,63 @@ | ||
package com.shinemo.mpush.core; | ||
|
||
|
||
import com.google.common.cache.Cache; | ||
import com.google.common.cache.CacheBuilder; | ||
import com.google.common.cache.RemovalListener; | ||
import com.google.common.cache.RemovalNotification; | ||
import com.shinemo.mpush.api.Connection; | ||
|
||
import java.util.Map; | ||
import java.util.concurrent.ConcurrentHashMap; | ||
import io.netty.util.internal.chmv8.ConcurrentHashMapV8; | ||
|
||
import java.util.concurrent.Callable; | ||
import java.util.concurrent.ConcurrentMap; | ||
import java.util.concurrent.ExecutionException; | ||
import java.util.concurrent.TimeUnit; | ||
|
||
/** | ||
* Created by ohun on 2015/12/22. | ||
*/ | ||
public class ConnectionManager { | ||
public static final ConnectionManager INSTANCE = new ConnectionManager(); | ||
private Map<String, Connection> connections = new ConcurrentHashMap<String, Connection>(); | ||
// private final ConcurrentMap<String, NettyConnection> connections = new ConcurrentHashMapV8<String, NettyConnection>(); | ||
|
||
public Connection get(String channelId) { | ||
return connections.get(channelId); | ||
private final Cache<String,NettyConnection> cacherClients = CacheBuilder.newBuilder() | ||
.maximumSize(2<<17) | ||
.expireAfterAccess(27, TimeUnit.MINUTES) | ||
.removalListener(new RemovalListener<String, NettyConnection>() { | ||
public void onRemoval(RemovalNotification<String,NettyConnection> notification) { | ||
if(notification.getValue().isClosed()){ | ||
// notification.getValue().close("[Remoting] removed from cache"); | ||
} | ||
}; | ||
}).build(); | ||
|
||
public Connection get(final String channelId) throws ExecutionException { | ||
|
||
NettyConnection client = cacherClients.get(channelId, new Callable<NettyConnection>() { | ||
@Override | ||
public NettyConnection call() throws Exception { | ||
NettyConnection client = getFromRedis(channelId); | ||
return client; | ||
} | ||
}); | ||
if (client == null || !client.isClosed()) { | ||
cacherClients.invalidate(channelId); | ||
return null; | ||
} | ||
return client; | ||
|
||
} | ||
|
||
public void add(Connection connection) { | ||
connections.put(connection.getId(), connection); | ||
public void add(NettyConnection connection) { | ||
cacherClients.put(connection.getId(), connection); | ||
} | ||
|
||
public void remove(Connection connection) { | ||
connections.remove(connection.getId()); | ||
public void remove(String channelId) { | ||
cacherClients.invalidate(channelId); | ||
} | ||
|
||
private NettyConnection getFromRedis(String channelId){ | ||
return null; | ||
} | ||
} |
Oops, something went wrong.