mirror of
https://github.com/ProjectSWGCore/Holocore.git
synced 2026-10-06 21:13:37 -04:00
Further improved memory performance of the TCPServer
This commit is contained in:
@@ -33,14 +33,15 @@ import java.net.InetSocketAddress;
|
||||
import java.net.Socket;
|
||||
import java.net.SocketAddress;
|
||||
import java.nio.ByteBuffer;
|
||||
import java.nio.channels.ClosedChannelException;
|
||||
import java.nio.channels.SelectableChannel;
|
||||
import java.nio.channels.SelectionKey;
|
||||
import java.nio.channels.Selector;
|
||||
import java.nio.channels.ServerSocketChannel;
|
||||
import java.nio.channels.SocketChannel;
|
||||
import java.util.HashMap;
|
||||
import java.util.Iterator;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
|
||||
|
||||
public class TCPServer {
|
||||
@@ -145,41 +146,6 @@ public class TCPServer {
|
||||
this.callback = callback;
|
||||
}
|
||||
|
||||
private void accept() {
|
||||
try {
|
||||
SocketChannel sc = channel.accept();
|
||||
if (sc == null)
|
||||
return;
|
||||
sc.configureBlocking(false);
|
||||
sockets.put(sc.getRemoteAddress(), sc);
|
||||
if (callback != null)
|
||||
callback.onIncomingConnection(sc.socket());
|
||||
} catch (IOException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
}
|
||||
|
||||
private void read(SocketChannel s) {
|
||||
ByteBuffer data = ByteBuffer.allocate(bufferSize);
|
||||
try {
|
||||
int n = s.read(data);
|
||||
if (n == -1) {
|
||||
disconnect(s);
|
||||
return;
|
||||
}
|
||||
if (n == 0)
|
||||
return;
|
||||
data.flip();
|
||||
ByteBuffer smaller = ByteBuffer.allocate(n);
|
||||
smaller.put(data.array(), 0, n);
|
||||
if (callback != null)
|
||||
callback.onIncomingData(s.socket(), smaller.array());
|
||||
} catch (IOException e) {
|
||||
e.printStackTrace();
|
||||
disconnect(s);
|
||||
}
|
||||
}
|
||||
|
||||
public interface TCPCallback {
|
||||
void onIncomingConnection(Socket s);
|
||||
void onConnectionDisconnect(Socket s);
|
||||
@@ -188,51 +154,41 @@ public class TCPServer {
|
||||
|
||||
private class TCPListener implements Runnable {
|
||||
|
||||
private final ByteBuffer buffer;
|
||||
private Thread thread;
|
||||
private boolean running;
|
||||
private boolean changedSelector;
|
||||
|
||||
public TCPListener() {
|
||||
buffer = ByteBuffer.allocateDirect(bufferSize);
|
||||
running = false;
|
||||
thread = null;
|
||||
changedSelector = false;
|
||||
}
|
||||
|
||||
public void start() {
|
||||
running = true;
|
||||
changedSelector = false;
|
||||
thread = new Thread(this);
|
||||
thread.start();
|
||||
}
|
||||
|
||||
public void stop() {
|
||||
running = false;
|
||||
changedSelector = false;
|
||||
if (thread != null)
|
||||
thread.interrupt();
|
||||
thread = null;
|
||||
}
|
||||
|
||||
public void run() {
|
||||
changedSelector = true;
|
||||
Selector selector = null;
|
||||
try {
|
||||
try (Selector selector = setupSelector()) {
|
||||
while (running) {
|
||||
try {
|
||||
if (changedSelector) {
|
||||
selector = setupSelector();
|
||||
changedSelector = false;
|
||||
}
|
||||
if (selector.select() > 0)
|
||||
processSelectionKeys(selector.selectedKeys());
|
||||
} catch (IOException e) {
|
||||
selector.select();
|
||||
processSelectionKeys(selector);
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
if (selector != null) {
|
||||
try { selector.close(); } catch (IOException e) { }
|
||||
}
|
||||
} catch (IOException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -248,11 +204,14 @@ public class TCPServer {
|
||||
return selector;
|
||||
}
|
||||
|
||||
private void processSelectionKeys(Set<SelectionKey> keys) {
|
||||
for (SelectionKey key : keys) {
|
||||
private void processSelectionKeys(Selector selector) throws ClosedChannelException {
|
||||
Iterator<SelectionKey> it = selector.selectedKeys().iterator();
|
||||
while (it.hasNext()) {
|
||||
SelectionKey key = it.next();
|
||||
if (key.isAcceptable()) {
|
||||
accept();
|
||||
changedSelector = true;
|
||||
SocketChannel sc = accept();
|
||||
if (sc != null)
|
||||
sc.register(selector, SelectionKey.OP_READ | SelectionKey.OP_CONNECT);
|
||||
} else if (key.isReadable()) {
|
||||
SelectableChannel selectable = key.channel();
|
||||
if (selectable instanceof SocketChannel)
|
||||
@@ -261,9 +220,47 @@ public class TCPServer {
|
||||
SelectableChannel selectable = key.channel();
|
||||
if (!selectable.isOpen() && selectable instanceof SocketChannel) {
|
||||
disconnect((SocketChannel) selectable);
|
||||
changedSelector = true;
|
||||
}
|
||||
}
|
||||
it.remove();
|
||||
}
|
||||
}
|
||||
|
||||
private SocketChannel accept() {
|
||||
SocketChannel sc = null;
|
||||
try {
|
||||
sc = channel.accept();
|
||||
if (sc == null)
|
||||
return sc;
|
||||
sc.configureBlocking(false);
|
||||
sockets.put(sc.getRemoteAddress(), sc);
|
||||
if (callback != null)
|
||||
callback.onIncomingConnection(sc.socket());
|
||||
} catch (IOException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
return sc;
|
||||
}
|
||||
|
||||
private void read(SocketChannel s) {
|
||||
try {
|
||||
buffer.position(0);
|
||||
buffer.limit(bufferSize);
|
||||
int n = s.read(buffer);
|
||||
buffer.flip();
|
||||
if (n == -1) {
|
||||
disconnect(s);
|
||||
return;
|
||||
}
|
||||
if (n == 0)
|
||||
return;
|
||||
ByteBuffer smaller = ByteBuffer.allocate(n);
|
||||
smaller.put(buffer);
|
||||
if (callback != null)
|
||||
callback.onIncomingData(s.socket(), smaller.array());
|
||||
} catch (IOException e) {
|
||||
e.printStackTrace();
|
||||
disconnect(s);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user