Various bugfixes and improvements

This commit is contained in:
Josh Larson
2018-08-30 15:57:18 -05:00
parent 3f980c0c32
commit 668025fac8
15 changed files with 499 additions and 441 deletions
@@ -7,6 +7,7 @@ import com.projectswg.forwarder.intents.control.StartForwarderIntent;
import com.projectswg.forwarder.intents.control.StopForwarderIntent;
import me.joshlarson.jlcommon.concurrency.Delay;
import me.joshlarson.jlcommon.control.IntentManager;
import me.joshlarson.jlcommon.control.IntentManager.IntentSpeedStatistics;
import me.joshlarson.jlcommon.control.Manager;
import me.joshlarson.jlcommon.control.SafeMain;
import me.joshlarson.jlcommon.log.Log;
@@ -19,6 +20,7 @@ import java.net.InetSocketAddress;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.util.Collections;
import java.util.Comparator;
import java.util.List;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.zip.ZipEntry;
@@ -97,9 +99,8 @@ public class Forwarder {
}
public void run() {
intentManager.initialize();
intentManager.registerForIntent(ClientConnectedIntent.class, cci -> connected.set(true));
intentManager.registerForIntent(ClientDisconnectedIntent.class, cdi -> connected.set(false));
intentManager.registerForIntent(ClientConnectedIntent.class, "Forwarder#handleClientConnectedIntent", cci -> connected.set(true));
intentManager.registerForIntent(ClientDisconnectedIntent.class, "Forwarder#handleClientDisconnectedIntent", cdi -> connected.set(false));
ConnectionManager primary = new ConnectionManager();
{
@@ -108,13 +109,27 @@ public class Forwarder {
Manager.start(managers);
new StartForwarderIntent(data).broadcast(intentManager);
Manager.run(managers, 100);
List<IntentSpeedStatistics> intentTimes = intentManager.getSpeedRecorder();
intentTimes.sort(Comparator.comparingLong(IntentSpeedStatistics::getTotalTime).reversed());
Log.i(" Intent Times: [%d]", intentTimes.size());
Log.i(" %-30s%-60s%-40s%-10s%-20s", "Intent", "Receiver Class", "Receiver Method", "Count", "Time");
for (IntentSpeedStatistics record : intentTimes) {
String receiverName = record.getKey().toString();
if (receiverName.indexOf('$') != -1)
receiverName = receiverName.substring(0, receiverName.indexOf('$'));
receiverName = receiverName.replace("com.projectswg.forwarder.services.", "");
String intentName = record.getIntent().getSimpleName();
String recordCount = Long.toString(record.getCount());
String recordTime = String.format("%.6fms", record.getTotalTime() / 1E6);
String [] receiverSplit = receiverName.split("#", 2);
Log.i(" %-30s%-60s%-40s%-10s%-20s", intentName, receiverSplit[0], receiverSplit[1], recordCount, recordTime);
}
new StopForwarderIntent().broadcast(intentManager);
Manager.stop(managers);
}
intentManager.terminate(false);
intentManager.close(false, 1000);
primary.setIntentManager(null);
}
public ForwarderData getData() {
@@ -0,0 +1,31 @@
/***********************************************************************************
* Copyright (C) 2018 /// Project SWG /// www.projectswg.com *
* *
* This file is part of the ProjectSWG Launcher. *
* *
* This program is free software: you can redistribute it and/or modify *
* it under the terms of the GNU Affero General Public License as published by *
* the Free Software Foundation, either version 3 of the License, or *
* (at your option) any later version. *
* *
* This program is distributed in the hope that it will be useful, *
* but WITHOUT ANY WARRANTY; without even the implied warranty of *
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the *
* GNU Affero General Public License for more details. *
* *
* You should have received a copy of the GNU Affero General Public License *
* along with this program. If not, see <https://www.gnu.org/licenses/>. *
* *
***********************************************************************************/
package com.projectswg.forwarder.intents.server;
import me.joshlarson.jlcommon.control.Intent;
public class RequestServerConnectionIntent extends Intent {
public RequestServerConnectionIntent() {
}
}
@@ -1,9 +1,8 @@
package com.projectswg.forwarder.resources.networking.data;
import com.projectswg.common.network.NetBuffer;
import com.projectswg.forwarder.resources.networking.packets.Fragmented;
import com.projectswg.forwarder.resources.networking.packets.Packet;
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.List;
@@ -22,27 +21,25 @@ public class FragmentedProcessor {
public byte [] addFragmented(Fragmented frag) {
fragmentedBuffer.add(frag);
frag = fragmentedBuffer.get(0);
ByteBuffer data = frag.getPacketData();
data.position(4);
int size = Packet.getNetInt(data);
int index = data.remaining();
for (int i = 1; i < fragmentedBuffer.size() && index < size; i++)
index += fragmentedBuffer.get(i).getPacketData().limit()-4;
if (index == size)
return processFragmentedReady(size);
return null;
NetBuffer payload = NetBuffer.wrap(frag.getPayload());
int length = payload.getNetInt();
if ((int) ((length+4.0) / 489) > fragmentedBuffer.size())
return null; // Doesn't have the minimum number of required packets
return processFragmentedReady(length);
}
private byte [] processFragmentedReady(int size) {
byte [] combined = new byte[size];
int index = 0;
while (index < combined.length) {
ByteBuffer packet = fragmentedBuffer.get(0).getPacketData();
packet.position(index == 0 ? 8 : 4);
int len = packet.remaining();
packet.get(combined, index, len);
index += len;
fragmentedBuffer.remove(0);
byte [] payload = fragmentedBuffer.remove(0).getPayload();
int header = (index == 0) ? 4 : 0;
System.arraycopy(payload, header, combined, index, payload.length - header);
index += payload.length - header;
}
return combined;
}
@@ -1,26 +1,26 @@
package com.projectswg.forwarder.resources.networking.data;
import com.projectswg.forwarder.resources.networking.data.ProtocolStack.ConnectionStream;
import com.projectswg.forwarder.resources.networking.packets.DataChannel;
import com.projectswg.forwarder.resources.networking.packets.Fragmented;
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.List;
import java.util.Queue;
import java.util.concurrent.atomic.AtomicInteger;
public class Packager {
private final AtomicInteger size;
private final DataChannel channel;
private final List<byte[]> dataChannel;
private final Queue<byte[]> outboundRaw;
private final Queue<SequencedOutbound> outboundPackaged;
private final ProtocolStack stack;
private final ConnectionStream<SequencedOutbound> outboundPackaged;
public Packager(Queue<byte[]> outboundRaw, Queue<SequencedOutbound> outboundPackaged, ProtocolStack stack) {
public Packager(Queue<byte[]> outboundRaw, ConnectionStream<SequencedOutbound> outboundPackaged, ProtocolStack stack) {
this.size = new AtomicInteger(8);
this.channel = new DataChannel();
this.dataChannel = new ArrayList<>();
this.outboundRaw = outboundRaw;
this.outboundPackaged = outboundPackaged;
this.stack = stack;
}
public void handle(int maxPackaged) {
@@ -46,29 +46,27 @@ public class Packager {
}
private void addToDataChannel(byte [] packet, int packetSize) {
channel.addPacket(packet);
dataChannel.add(packet);
size.getAndAdd(packetSize);
}
private void sendDataChannel() {
if (channel.getPacketCount() == 0)
if (dataChannel.isEmpty())
return;
channel.setSequence(stack.getAndIncrementTxSequence());
outboundPackaged.add(new SequencedOutbound(channel.getSequence(), channel.encode().array()));
outboundPackaged.addOrdered(new SequencedOutbound(new DataChannel(dataChannel)));
reset();
}
private void sendFragmented(byte [] packet) {
Fragmented[] frags = Fragmented.encode(ByteBuffer.wrap(packet), stack.getTxSequence());
stack.getAndIncrementTxSequence(frags.length);
for (Fragmented frag : frags) {
outboundPackaged.add(new SequencedOutbound(frag.getSequence(), frag.encode().array()));
byte[][] frags = Fragmented.split(packet);
for (byte [] frag : frags) {
outboundPackaged.addOrdered(new SequencedOutbound(new Fragmented((short) 0, frag)));
}
}
private void reset() {
channel.clearPackets();
dataChannel.clear();
size.set(8);
}
@@ -13,37 +13,29 @@ import java.util.function.BiConsumer;
public class ProtocolStack {
private final PriorityQueue<SequencedPacket> sequenced;
private final FragmentedProcessor fragmentedProcessor;
private final InetSocketAddress source;
private final BiConsumer<InetSocketAddress, byte[]> sender;
private final ClientServer server;
private final Queue<byte []> outboundRaw;
private final Queue<SequencedOutbound> outboundPackaged;
private final ConnectionStream<SequencedPacket> inbound;
private final ConnectionStream<SequencedOutbound> outbound;
private final Packager packager;
private final Object txMutex;
private InetSocketAddress pingSource;
private int connectionId;
private short rxSequence;
private short txSequence;
private boolean txOverflow;
public ProtocolStack(InetSocketAddress source, ClientServer server, BiConsumer<InetSocketAddress, byte[]> sender) {
this.sequenced = new PriorityQueue<>();
this.fragmentedProcessor = new FragmentedProcessor();
this.source = source;
this.sender = sender;
this.server = server;
this.outboundRaw = new LinkedList<>();
this.outboundPackaged = new LinkedList<>();
this.packager = new Packager(outboundRaw, outboundPackaged, this);
this.txMutex = new Object();
this.inbound = new ConnectionStream<>();
this.outbound = new ConnectionStream<>();
this.packager = new Packager(outboundRaw, outbound, this);
this.connectionId = 0;
this.rxSequence = 0;
this.txSequence = 0;
this.txOverflow = false;
}
public void send(Packet packet) {
@@ -74,11 +66,11 @@ public class ProtocolStack {
}
public short getRxSequence() {
return rxSequence;
return inbound.getSequence();
}
public short getTxSequence() {
return txSequence;
return outbound.getSequence();
}
public void setPingSource(InetSocketAddress source) {
@@ -89,46 +81,12 @@ public class ProtocolStack {
this.connectionId = connectionId;
}
public short getAndIncrementTxSequence() {
return getAndIncrementTxSequence(1);
}
public short getAndIncrementTxSequence(int amount) {
synchronized (txMutex) {
short prev = this.txSequence;
short next = prev;
next += amount;
if (prev > next)
txOverflow = true;
this.txSequence = next;
return prev;
}
}
public boolean addIncoming(@NotNull SequencedPacket packet) {
synchronized (sequenced) {
if (packet.getSequence() < rxSequence)
return true;
// If it already exists in here, don't add it again
for (SequencedPacket seq : sequenced) {
if (seq.getSequence() == packet.getSequence())
return true;
}
sequenced.add(packet);
packet = sequenced.peek();
assert packet != null : "the world is on fire";
return packet.getSequence() == rxSequence;
}
public SequencedStatus addIncoming(@NotNull SequencedPacket packet) {
return inbound.addUnordered(packet);
}
public SequencedPacket getNextIncoming() {
synchronized (sequenced) {
SequencedPacket peek = sequenced.peek();
if (peek == null || peek.getSequence() != rxSequence)
return null;
rxSequence++;
return sequenced.poll();
}
return inbound.poll();
}
public byte [] addFragmented(Fragmented frag) {
@@ -140,32 +98,22 @@ public class ProtocolStack {
}
public short getFirstUnacknowledgedOutbound() {
SequencedOutbound out = outboundPackaged.peek();
SequencedOutbound out = outbound.peek();
if (out == null)
return -1;
return out.getSequence();
}
public void clearAcknowledgedOutbound(short sequence) {
synchronized (txMutex) {
if (txOverflow && sequence <= txSequence) {
outboundPackaged.removeIf(out -> out.getSequence() > txSequence);
txOverflow = false;
}
SequencedOutbound out = outboundPackaged.peek();
while (out != null && out.getSequence() <= sequence) {
outboundPackaged.poll();
out = outboundPackaged.peek();
}
}
outbound.removeOrdered(sequence);
}
public void fillOutboundPackagedBuffer(int maxPackaged) {
packager.handle(maxPackaged);
}
public Collection<SequencedOutbound> getOutboundPackagedBuffer() {
return Collections.unmodifiableCollection(outboundPackaged);
public int fillOutboundBuffer(SequencedOutbound [] buffer) {
return outbound.fillBuffer(buffer);
}
@Override
@@ -173,4 +121,92 @@ public class ProtocolStack {
return String.format("ProtocolStack[server=%s, source=%s, connectionId=%d]", server, source, connectionId);
}
public static class ConnectionStream<T extends SequencedPacket> {
private final PriorityQueue<T> sequenced;
private final PriorityQueue<T> queued;
private short sequence;
public ConnectionStream() {
this.sequenced = new PriorityQueue<>();
this.queued = new PriorityQueue<>();
this.sequence = 0;
}
public short getSequence() {
return sequence;
}
public synchronized SequencedStatus addUnordered(@NotNull T packet) {
if (SequencedPacket.compare(sequence, packet.getSequence()) > 0) {
T peek = peek();
return peek != null && peek.getSequence() == sequence ? SequencedStatus.READY : SequencedStatus.STALE;
}
if (packet.getSequence() == sequence) {
sequenced.add(packet);
sequence++;
// Add queued OOO packets
T queue;
while ((queue = queued.peek()) != null && queue.getSequence() == sequence) {
sequenced.add(queued.poll());
sequence++;
}
return SequencedStatus.READY;
} else {
queued.add(packet);
return SequencedStatus.OUT_OF_ORDER;
}
}
public synchronized void addOrdered(@NotNull T packet) {
packet.setSequence(sequence);
addUnordered(packet);
}
public synchronized void removeOrdered(short sequence) {
T packet;
List<Short> sequencesRemoved = new ArrayList<>();
while ((packet = sequenced.peek()) != null && SequencedPacket.compare(sequence, packet.getSequence()) >= 0) {
T removed = sequenced.poll();
assert packet == removed;
sequencesRemoved.add(packet.getSequence());
}
Log.t("Removed acknowledged: %s", sequencesRemoved);
}
public synchronized T peek() {
return sequenced.peek();
}
public synchronized T poll() {
return sequenced.poll();
}
public synchronized int fillBuffer(T [] buffer) {
int n = 0;
for (T packet : sequenced) {
if (n >= buffer.length)
break;
buffer[n++] = packet;
}
return n;
}
public int size() {
return sequenced.size();
}
}
public enum SequencedStatus {
READY,
OUT_OF_ORDER,
STALE
}
}
@@ -1,19 +1,24 @@
package com.projectswg.forwarder.resources.networking.data;
public class SequencedOutbound {
import com.projectswg.forwarder.resources.networking.packets.SequencedPacket;
import java.nio.ByteBuffer;
public class SequencedOutbound implements SequencedPacket {
private final short sequence;
private final byte [] data;
private final SequencedPacket packet;
private byte [] data;
private boolean sent;
public SequencedOutbound(short sequence, byte [] data) {
this.sequence = sequence;
this.data = data;
public SequencedOutbound(SequencedPacket packet) {
this.packet = packet;
this.data = packet.encode().array();
this.sent = false;
}
@Override
public short getSequence() {
return sequence;
return packet.getSequence();
}
public byte[] getData() {
@@ -24,6 +29,17 @@ public class SequencedOutbound {
return sent;
}
@Override
public ByteBuffer encode() {
return packet.encode();
}
@Override
public void setSequence(short sequence) {
this.packet.setSequence(sequence);
this.data = packet.encode().array();
}
public void setSent(boolean sent) {
this.sent = sent;
}
@@ -31,6 +31,7 @@ import com.projectswg.common.network.NetBuffer;
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
public class DataChannel extends Packet implements SequencedPacket {
@@ -48,6 +49,13 @@ public class DataChannel extends Packet implements SequencedPacket {
this.multiPacket = 0;
}
public DataChannel(List<byte[]> content) {
this.content = new ArrayList<>(content);
this.channel = Channel.DATA_CHANNEL_A;
this.sequence = 0;
this.multiPacket = (short) (content.isEmpty() ? 0 : 0x19);
}
public DataChannel(ByteBuffer data) {
this();
decode(data);
@@ -55,9 +63,7 @@ public class DataChannel extends Packet implements SequencedPacket {
public DataChannel(byte[][] packets) {
this();
for (byte[] p : packets) {
content.add(p);
}
content.addAll(Arrays.asList(packets));
}
@Override
@@ -74,7 +80,7 @@ public class DataChannel extends Packet implements SequencedPacket {
sequence = getNetShort(data);
multiPacket = getNetShort(data);
if (multiPacket == 0x19) {
int length = 0;
int length;
while (data.remaining() > 1) {
length = data.get() & 0xFF;
if (length == 0xFF)
@@ -97,11 +103,6 @@ public class DataChannel extends Packet implements SequencedPacket {
@Override
public ByteBuffer encode() {
return encode(this.sequence);
}
public ByteBuffer encode(int sequence) {
this.sequence = (short) sequence;
NetBuffer data;
if (content.size() == 1) {
byte[] pData = content.get(0);
@@ -150,20 +151,9 @@ public class DataChannel extends Packet implements SequencedPacket {
}
}
@Override
public int compareTo(SequencedPacket p) {
if (sequence < p.getSequence())
return -1;
if (sequence == p.getSequence())
return 0;
return 1;
}
@Override
public boolean equals(Object o) {
if (!(o instanceof DataChannel))
return false;
return ((DataChannel) o).sequence == sequence;
return o instanceof DataChannel && sequence == ((DataChannel) o).sequence;
}
@Override
@@ -27,71 +27,53 @@
***********************************************************************************/
package com.projectswg.forwarder.resources.networking.packets;
import com.projectswg.common.network.NetBuffer;
import java.nio.ByteBuffer;
import java.nio.ByteOrder;
public class Fragmented extends Packet implements SequencedPacket {
private short sequence;
private int length;
private ByteBuffer packet;
private ByteBuffer data;
private byte [] payload;
public Fragmented() {
this.sequence = 0;
this.length = 0;
this.packet = null;
this.data = null;
this.payload = null;
}
public Fragmented(ByteBuffer data) {
decode(data);
length = -1;
}
public Fragmented(ByteBuffer data, int sequence) {
this.sequence = (short) sequence;
decode(data);
length = -1;
}
public void setPacket(ByteBuffer packet) {
this.packet = packet;
}
public void decode(ByteBuffer data) {
data.position(2);
sequence = getNetShort(data);
this.data = data;
}
public ByteBuffer encode() {
return data;
}
public Fragmented [] encode(int startSequence) {
packet.position(0);
int ord = 0;
Fragmented [] packets = new Fragmented[(int) Math.ceil((packet.remaining()+4)/489.0)];
while (packet.remaining() > 0) {
packets[ord] = createSegment(startSequence++, ord++, packet);
}
return packets;
public Fragmented(short sequence, byte [] payload) {
this.sequence = sequence;
this.payload = payload;
}
@Override
public int compareTo(SequencedPacket p) {
if (sequence < p.getSequence())
return -1;
if (sequence == p.getSequence())
return 0;
return 1;
public void decode(ByteBuffer bb) {
super.decode(bb);
NetBuffer data = NetBuffer.wrap(bb);
short type = data.getNetShort();
assert (type - 0x0D) < 4;
this.sequence = data.getNetShort();
this.payload = data.getArray(data.remaining());
}
@Override
public ByteBuffer encode() {
NetBuffer data = NetBuffer.allocate(4 + payload.length);
data.addNetShort(0x0D);
data.addNetShort(sequence);
data.addRawArray(payload);
return data.getBuffer();
}
@Override
public boolean equals(Object o) {
if (!(o instanceof Fragmented))
return super.equals(o);
return ((Fragmented) o).sequence == sequence;
return o instanceof Fragmented && sequence == ((Fragmented) o).sequence;
}
@Override
@@ -99,43 +81,44 @@ public class Fragmented extends Packet implements SequencedPacket {
return sequence;
}
public ByteBuffer getPacketData() { return data; }
public short getSequence() { return sequence; }
public int getDatLength() { return length; }
@Override
public String toString() {
return String.format("Fragmented[seq=%d, len=%d]", sequence, payload.length);
}
public static final Fragmented [] encode(ByteBuffer data, int startSequence) {
data.position(0);
int ord = 0;
Fragmented [] packets = new Fragmented[(int) Math.ceil((data.remaining()+4)/489.0)];
while (data.remaining() > 0) {
packets[ord] = createSegment(startSequence++, ord++, data);
@Override
public short getSequence() {
return sequence;
}
public byte[] getPayload() {
return payload;
}
@Override
public void setSequence(short sequence) {
this.sequence = sequence;
}
public void setPayload(byte[] payload) {
this.payload = payload;
}
public static byte [][] split(byte [] data) {
int offset = 0;
int packetCount = (int) Math.ceil((data.length+4)/489.0);
byte [][] packets = new byte[packetCount][];
for (int i = 0; i < packetCount; i++) {
int header = (i == 0) ? 4 : 0;
ByteBuffer segment = ByteBuffer.allocate(Math.min(data.length-offset-header, 489)).order(ByteOrder.BIG_ENDIAN);
if (i == 0)
segment.putInt(data.length);
int segmentLength = segment.remaining();
segment.put(data, offset, segmentLength);
offset += segmentLength;
packets[i] = segment.array();
}
return packets;
}
private static final Fragmented createSegment(int startSequence, int ord, ByteBuffer packet) {
int header = (ord == 0) ? 8 : 4;
ByteBuffer data = ByteBuffer.allocate(Math.min(packet.remaining()+header, 493));
addNetShort(data, 0x0D);
addNetShort(data, startSequence);
if (ord == 0)
addNetInt(data, packet.remaining());
int len = data.remaining();
data.put(packet.array(), packet.position(), len);
packet.position(packet.position() + len);
byte [] pData = new byte[data.array().length-header];
System.arraycopy(data.array(), header, pData, 0, pData.length);
Fragmented f = new Fragmented(data, startSequence);
f.packet = ByteBuffer.wrap(pData);
return f;
}
@Override
public String toString() {
return String.format("Fragmented[seq=%d, len=%d]", sequence, length);
}
}
@@ -1,5 +1,34 @@
package com.projectswg.forwarder.resources.networking.packets;
import org.jetbrains.annotations.NotNull;
import java.nio.ByteBuffer;
public interface SequencedPacket extends Comparable<SequencedPacket> {
void setSequence(short sequence);
short getSequence();
ByteBuffer encode();
default int compareTo(@NotNull SequencedPacket p) {
return compare(this, p);
}
static int compare(@NotNull SequencedPacket a, @NotNull SequencedPacket b) {
return compare(a.getSequence(), b.getSequence());
}
static int compare(short a, short b) {
int aSeq = a & 0xFFFF;
int bSeq = b & 0xFFFF;
if (aSeq == bSeq)
return 0;
int diff = bSeq - aSeq;
if (diff <= 0)
diff += 0x10000;
return diff < 30000 ? -1 : 1;
}
}
@@ -43,30 +43,42 @@ public class ClientInboundDataService extends Service {
@Multiplexer
private void handleDataChannel(ProtocolStack stack, DataChannel data) {
if (stack.addIncoming(data)) {
readAvailablePackets(stack);
} else {
Log.d("Inbound Out of Order %d (data)", data.getSequence());
for (short seq = stack.getRxSequence(); seq < data.getSequence(); seq++)
stack.send(new OutOfOrder(seq));
switch (stack.addIncoming(data)) {
case READY:
readAvailablePackets(stack);
break;
case OUT_OF_ORDER:
Log.d("Data/Inbound Received Data Out of Order %d", data.getSequence());
for (short seq = stack.getRxSequence(); seq < data.getSequence(); seq++)
stack.send(new OutOfOrder(seq));
break;
default:
break;
}
}
@Multiplexer
private void handleFragmented(ProtocolStack stack, Fragmented frag) {
if (stack.addIncoming(frag)) {
readAvailablePackets(stack);
} else {
Log.d("Inbound Out of Order %d (frag)", frag.getSequence());
for (short seq = stack.getRxSequence(); seq < frag.getSequence(); seq++)
stack.send(new OutOfOrder(seq));
switch (stack.addIncoming(frag)) {
case READY:
readAvailablePackets(stack);
break;
case OUT_OF_ORDER:
Log.d("Data/Inbound Received Frag Out of Order %d", frag.getSequence());
for (short seq = stack.getRxSequence(); seq < frag.getSequence(); seq++)
stack.send(new OutOfOrder(seq));
break;
default:
break;
}
}
private void readAvailablePackets(ProtocolStack stack) {
short highestSequence = -1;
boolean updatedSequence = false;
SequencedPacket packet = stack.getNextIncoming();
while (packet != null) {
Log.t("Data/Inbound Received: %s", packet);
if (packet instanceof DataChannel) {
for (byte [] data : ((DataChannel) packet).getPackets())
onData(data);
@@ -75,19 +87,17 @@ public class ClientInboundDataService extends Service {
if (data != null)
onData(data);
}
Log.t("Data Inbound: %s", packet);
highestSequence = packet.getSequence();
packet = stack.getNextIncoming();
updatedSequence = true;
}
if (highestSequence != -1) {
Log.t("Inbound Acknowledge %d", highestSequence);
if (updatedSequence)
stack.send(new Acknowledge(highestSequence));
}
}
private void onData(byte [] data) {
PacketType type = PacketType.fromCrc(ByteBuffer.wrap(data).order(ByteOrder.LITTLE_ENDIAN).getInt(2));
Log.d("Incoming Data: %s", type);
Log.d("Data/Inbound Received Data: %s", type);
intentChain.broadcastAfter(getIntentManager(), new DataPacketInboundIntent(data));
}
@@ -12,8 +12,9 @@ import com.projectswg.forwarder.resources.networking.data.SequencedOutbound;
import com.projectswg.forwarder.resources.networking.packets.Acknowledge;
import com.projectswg.forwarder.resources.networking.packets.OutOfOrder;
import com.projectswg.forwarder.resources.networking.packets.Packet;
import me.joshlarson.jlcommon.concurrency.BasicThread;
import me.joshlarson.jlcommon.concurrency.Delay;
import me.joshlarson.jlcommon.concurrency.ScheduledThreadPool;
import me.joshlarson.jlcommon.concurrency.SmartLock;
import me.joshlarson.jlcommon.control.IntentHandler;
import me.joshlarson.jlcommon.control.IntentMultiplexer;
import me.joshlarson.jlcommon.control.IntentMultiplexer.Multiplexer;
@@ -24,28 +25,29 @@ import java.nio.ByteBuffer;
import java.nio.ByteOrder;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
public class ClientOutboundDataService extends Service {
private final IntentMultiplexer multiplexer;
private final Set<ProtocolStack> activeStacks;
private final ScheduledThreadPool timerThread;
private final Object outboundMutex;
private final BasicThread sendThread;
private final SmartLock signaller;
private ForwarderData data;
public ClientOutboundDataService() {
this.multiplexer = new IntentMultiplexer(this, ProtocolStack.class, Packet.class);
this.activeStacks = ConcurrentHashMap.newKeySet();
this.timerThread = new ScheduledThreadPool(2, 5, "outbound-sender-%d");
this.outboundMutex = new Object();
this.sendThread = new BasicThread("outbound-sender", this::persistentSend);
this.signaller = new SmartLock();
this.data = null;
}
@Override
public boolean terminate() {
timerThread.stop();
return timerThread.awaitTermination(1000);
sendThread.stop(true);
return sendThread.awaitTermination(1000);
}
@IntentHandler
@@ -60,22 +62,17 @@ public class ClientOutboundDataService extends Service {
@IntentHandler
private void handleClientConnectedIntent(ClientConnectedIntent cci) {
if (timerThread.isRunning())
if (sendThread.isExecuting())
return;
int interval = data.getOutboundTunerInterval();
if (interval <= 0)
interval = 20;
timerThread.start();
timerThread.executeWithFixedDelay(interval, interval, this::timerCallback);
timerThread.executeWithFixedRate(0, 5000, this::clearSentBit);
sendThread.start();
}
@IntentHandler
private void handleClientDisconnectedIntent(ClientDisconnectedIntent cdi) {
if (!timerThread.isRunning())
if (!sendThread.isExecuting())
return;
timerThread.stop();
timerThread.awaitTermination(500);
sendThread.stop(true);
sendThread.awaitTermination(1000);
}
@IntentHandler
@@ -92,103 +89,101 @@ public class ClientOutboundDataService extends Service {
private void handleDataPacketOutboundIntent(DataPacketOutboundIntent dpoi) {
PacketType type = PacketType.fromCrc(ByteBuffer.wrap(dpoi.getData()).order(ByteOrder.LITTLE_ENDIAN).getInt(2));
ClientServer filterServer = ClientServer.ZONE;
switch (type) {
case ERROR_MESSAGE:
case SERVER_ID:
case SERVER_NOW_EPOCH_TIME:
filterServer = null;
break;
case LOGIN_CLUSTER_STATUS:
case LOGIN_CLIENT_TOKEN:
case LOGIN_INCORRECT_CLIENT_ID:
case LOGIN_ENUM_CLUSTER:
case ENUMERATE_CHARACTER_ID:
case CHARACTER_CREATION_DISABLED:
filterServer = ClientServer.LOGIN;
break;
case HEART_BEAT_MESSAGE: {
HeartBeat heartbeat = new HeartBeat();
heartbeat.decode(NetBuffer.wrap(dpoi.getData()));
if (heartbeat.getPayload().length > 0) {
for (ProtocolStack stack : activeStacks)
stack.sendPing(heartbeat.getPayload());
return;
}
break;
}
}
synchronized (outboundMutex) {
if (filterServer == null) {
Log.d("Sending %d bytes to %s", dpoi.getData().length, activeStacks);
for (ProtocolStack stack : activeStacks)
stack.addOutbound(dpoi.getData());
} else {
final ClientServer finalFilterServer = filterServer;
ProtocolStack stack = activeStacks.stream().filter(s -> s.getServer() == finalFilterServer).findFirst().orElse(null);
if (stack != null) {
Log.d("Sending %d bytes to %s", dpoi.getData().length, stack);
stack.addOutbound(dpoi.getData());
if (type != null) {
switch (type) {
case ERROR_MESSAGE:
case SERVER_ID:
case SERVER_NOW_EPOCH_TIME:
filterServer = null;
break;
case LOGIN_CLUSTER_STATUS:
case LOGIN_CLIENT_TOKEN:
case LOGIN_INCORRECT_CLIENT_ID:
case LOGIN_ENUM_CLUSTER:
case ENUMERATE_CHARACTER_ID:
case CHARACTER_CREATION_DISABLED:
case DELETE_CHARACTER_REQUEST:
case DELETE_CHARACTER_RESPONSE:
filterServer = ClientServer.LOGIN;
break;
case HEART_BEAT_MESSAGE: {
HeartBeat heartbeat = new HeartBeat();
heartbeat.decode(NetBuffer.wrap(dpoi.getData()));
if (heartbeat.getPayload().length > 0) {
for (ProtocolStack stack : activeStacks)
stack.sendPing(heartbeat.getPayload());
return;
}
break;
}
}
}
final ClientServer finalFilterServer = filterServer;
ProtocolStack stack = activeStacks.stream().filter(s -> s.getServer() == finalFilterServer).findFirst().orElse(null);
if (stack == null) {
Log.d("Data/Oubound Sending %s [len=%d] to %s", type, dpoi.getData().length, activeStacks);
for (ProtocolStack active : activeStacks)
active.addOutbound(dpoi.getData());
} else {
Log.d("Data/Outbound Sending %s [len=%d] to %s", type, dpoi.getData().length, stack);
stack.addOutbound(dpoi.getData());
}
signaller.signal();
}
@Multiplexer
private void handleAcknowledgement(ProtocolStack stack, Acknowledge ack) {
Log.t("Acknowledged: %d. Min Sequence: %d", ack.getSequence(), stack.getFirstUnacknowledgedOutbound());
synchronized (outboundMutex) {
stack.clearAcknowledgedOutbound(ack.getSequence());
for (SequencedOutbound outbound : stack.getOutboundPackagedBuffer()) {
outbound.setSent(false);
}
}
Log.t("Data/Outbound Client Acknowledged: %d. Min Sequence: %d", ack.getSequence(), stack.getFirstUnacknowledgedOutbound());
stack.clearAcknowledgedOutbound(ack.getSequence());
}
@Multiplexer
private void handleOutOfOrder(ProtocolStack stack, OutOfOrder ooo) {
synchronized (outboundMutex) {
for (SequencedOutbound outbound : stack.getOutboundPackagedBuffer()) {
if (outbound.getSequence() > ooo.getSequence())
break;
outbound.setSent(false);
}
}
Log.t("Data/Outbound Out of Order: %d. Min Sequence: %d", ooo.getSequence(), stack.getFirstUnacknowledgedOutbound());
}
private void clearSentBit() {
synchronized (outboundMutex) {
for (ProtocolStack stack : activeStacks) {
for (SequencedOutbound outbound : stack.getOutboundPackagedBuffer()) {
outbound.setSent(false);
}
}
}
}
private void timerCallback() {
private void persistentSend() {
int maxSend = data.getOutboundTunerMaxSend();
if (maxSend <= 0)
maxSend = 100;
synchronized (outboundMutex) {
for (ProtocolStack stack : activeStacks) {
stack.fillOutboundPackagedBuffer(maxSend);
int sent = 0;
int runStart = Integer.MIN_VALUE;
int runEnd = 0;
for (SequencedOutbound outbound : stack.getOutboundPackagedBuffer()) {
runEnd = outbound.getSequence();
if (runStart == Integer.MIN_VALUE) {
runStart = runEnd;
maxSend = 400;
int interval = data.getOutboundTunerInterval();
if (interval <= 0)
interval = 1000;
int bucket = 0;
int [] sendCounts = new int[interval / 20];
SequencedOutbound [] buffer = new SequencedOutbound[maxSend];
Log.d("Data/Outbound Starting Persistent Send");
while (!Delay.isInterrupted()) {
sendCounts[bucket] = 0;
int previousInterval = 0;
for (int count : sendCounts)
previousInterval += count;
if (previousInterval < maxSend) {
for (ProtocolStack stack : activeStacks) {
stack.fillOutboundPackagedBuffer(maxSend);
int count = stack.fillOutboundBuffer(buffer);
for (int i = 0; i < count && Delay.sleepMicro(50); i++) {
stack.send(buffer[i].getData());
}
stack.send(outbound.getData());
Delay.sleepMicro(50);
if (sent++ >= maxSend)
break;
if (count > 0)
Log.t("Data/Outbound Sent %d Start: %d", count, buffer[0].getSequence());
sendCounts[bucket] += count;
}
if (runStart != Integer.MIN_VALUE)
Log.t("Sending to %s: %d - %d", stack.getSource(), runStart, runEnd);
try {
signaller.await(20, TimeUnit.MILLISECONDS);
} catch (InterruptedException e) {
break;
}
} else {
Delay.sleepMilli(20);
}
bucket = (bucket + 1) % sendCounts.length;
}
Log.d("Data/Outbound Stopping Persistent Send");
}
}
@@ -4,6 +4,8 @@ import com.projectswg.forwarder.Forwarder.ForwarderData;
import com.projectswg.forwarder.intents.client.*;
import com.projectswg.forwarder.intents.control.StartForwarderIntent;
import com.projectswg.forwarder.intents.control.StopForwarderIntent;
import com.projectswg.forwarder.intents.server.RequestServerConnectionIntent;
import com.projectswg.forwarder.intents.server.ServerConnectedIntent;
import com.projectswg.forwarder.intents.server.ServerDisconnectedIntent;
import com.projectswg.forwarder.resources.networking.ClientServer;
import com.projectswg.forwarder.resources.networking.data.ProtocolStack;
@@ -20,11 +22,15 @@ import java.net.*;
import java.nio.ByteBuffer;
import java.util.EnumMap;
import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
public class ClientServerService extends Service {
private final IntentChain intentChain;
private final Map<ClientServer, ProtocolStack> stacks;
private final AtomicBoolean serverConnection;
private final AtomicReference<SessionRequest> cachedSessionRequest;
private ForwarderData data;
private UDPServer loginServer;
@@ -33,6 +39,8 @@ public class ClientServerService extends Service {
public ClientServerService() {
this.intentChain = new IntentChain();
this.stacks = new EnumMap<>(ClientServer.class);
this.serverConnection = new AtomicBoolean(false);
this.cachedSessionRequest = new AtomicReference<>(null);
this.data = null;
this.loginServer = null;
this.zoneServer = null;
@@ -93,8 +101,19 @@ public class ClientServerService extends Service {
Log.i("Closed the login and zone udp servers");
}
@IntentHandler
private void handleServerConnectedIntent(ServerConnectedIntent sci) {
serverConnection.set(true);
SessionRequest request = this.cachedSessionRequest.getAndSet(null);
if (request != null) {
byte [] data = request.encode().array();
onLoginPacket(new DatagramPacket(data, data.length, new InetSocketAddress(request.getAddress(), request.getPort())));
}
}
@IntentHandler
private void handleServerDisconnectedIntent(ServerDisconnectedIntent sdi) {
serverConnection.set(false);
closeConnection(ClientServer.LOGIN);
closeConnection(ClientServer.ZONE);
}
@@ -102,10 +121,10 @@ public class ClientServerService extends Service {
private void customizeUdpServer(DatagramSocket socket) {
try {
socket.setReuseAddress(false);
socket.setTrafficClass(0x02 | 0x04 | 0x08 | 0x10);
socket.setTrafficClass(0x10);
socket.setBroadcast(false);
socket.setReceiveBufferSize(496 * 2048);
socket.setSendBufferSize(496 * 2048);
socket.setReceiveBufferSize(64 * 1024);
socket.setSendBufferSize(64 * 1024);
} catch (SocketException e) {
Log.w(e);
}
@@ -134,6 +153,8 @@ public class ClientServerService extends Service {
Packet parsed = parse(data);
if (parsed == null)
return;
parsed.setAddress(source.getAddress());
parsed.setPort(source.getPort());
if (parsed instanceof MultiPacket) {
for (byte [] child : ((MultiPacket) parsed).getPackets()) {
process(source, server, child);
@@ -154,6 +175,13 @@ public class ClientServerService extends Service {
}
private ProtocolStack process(InetSocketAddress source, ClientServer server, Packet parsed) {
if (!serverConnection.get()) {
if (parsed instanceof SessionRequest) {
cachedSessionRequest.set((SessionRequest) parsed);
intentChain.broadcastAfter(getIntentManager(), new RequestServerConnectionIntent());
}
return null;
}
if (parsed instanceof SessionRequest)
return onSessionRequest(source, server, (SessionRequest) parsed);
if (parsed instanceof Disconnect)
@@ -173,6 +201,11 @@ public class ClientServerService extends Service {
}
private ProtocolStack onSessionRequest(InetSocketAddress source, ClientServer server, SessionRequest request) {
ProtocolStack current = stacks.get(server);
if (current != null && request.getConnectionId() != current.getConnectionId()) {
closeConnection(server);
return null;
}
ProtocolStack stack = new ProtocolStack(source, server, (remote, data) -> send(remote, server, data));
stack.setConnectionId(request.getConnectionId());
@@ -1,40 +1,29 @@
package com.projectswg.forwarder.services.server;
import com.projectswg.common.utilities.ByteUtilities;
import com.projectswg.connection.HolocoreSocket;
import com.projectswg.connection.RawPacket;
import com.projectswg.connection.ServerConnectionChangedReason;
import com.projectswg.forwarder.Forwarder.ForwarderData;
import com.projectswg.forwarder.intents.client.ClientConnectedIntent;
import com.projectswg.forwarder.intents.client.ClientDisconnectedIntent;
import com.projectswg.forwarder.intents.client.DataPacketInboundIntent;
import com.projectswg.forwarder.intents.client.DataPacketOutboundIntent;
import com.projectswg.forwarder.intents.control.StartForwarderIntent;
import com.projectswg.forwarder.intents.control.StopForwarderIntent;
import com.projectswg.forwarder.intents.server.RequestServerConnectionIntent;
import com.projectswg.forwarder.intents.server.ServerConnectedIntent;
import com.projectswg.forwarder.intents.server.ServerDisconnectedIntent;
import com.projectswg.forwarder.resources.networking.NetInterceptor;
import me.joshlarson.jlcommon.concurrency.Delay;
import me.joshlarson.jlcommon.concurrency.ThreadPool;
import me.joshlarson.jlcommon.concurrency.BasicThread;
import me.joshlarson.jlcommon.control.IntentChain;
import me.joshlarson.jlcommon.control.IntentHandler;
import me.joshlarson.jlcommon.control.Service;
import me.joshlarson.jlcommon.log.Log;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
public class ServerConnectionService extends Service {
private static final int CONNECT_TIMEOUT = 5000;
private final IntentChain intentChain;
private final AtomicReference<Thread> activeThread;
private final AtomicBoolean running;
private final ThreadPool thread;
private final Lock runLock;
private final Lock sleepLock;
private final BasicThread thread;
private HolocoreSocket holocore;
private NetInterceptor interceptor;
@@ -42,27 +31,15 @@ public class ServerConnectionService extends Service {
public ServerConnectionService() {
this.intentChain = new IntentChain();
this.activeThread = new AtomicReference<>(null);
this.running = new AtomicBoolean(false);
this.thread = new ThreadPool(2, "server-connection");
this.runLock = new ReentrantLock(true);
this.sleepLock = new ReentrantLock(true);
this.thread = new BasicThread("server-connection", this::primaryConnectionLoop);
this.holocore = null;
this.interceptor = null;
this.data = null;
}
@Override
public boolean start() {
thread.start();
return true;
}
@Override
public boolean stop() {
running.set(false);
thread.stop(true);
return thread.awaitTermination(1000);
return stopRunningLoop();
}
@IntentHandler
@@ -77,8 +54,10 @@ public class ServerConnectionService extends Service {
}
@IntentHandler
private void handleClientConnectedIntent(ClientConnectedIntent cci) {
queueConnectionLoop(0);
private void handleRequestServerConnectionIntent(RequestServerConnectionIntent rsci) {
if (thread.isExecuting())
stopRunningLoop();
thread.start();
}
@IntentHandler
@@ -88,97 +67,43 @@ public class ServerConnectionService extends Service {
@IntentHandler
private void handleDataPacketInboundIntent(DataPacketInboundIntent dpii) {
if (running.get())
HolocoreSocket holocore = this.holocore;
if (holocore != null)
holocore.send(interceptor.interceptClient(dpii.getData()));
else
Log.w("Dropping packet destined for server: %s", ByteUtilities.getHexString(dpii.getData()));
}
private void stopRunningLoop() {
running.set(false);
Thread activeThread = this.activeThread.get();
if (activeThread != null)
activeThread.interrupt();
try {
if (runLock.tryLock(1, TimeUnit.SECONDS))
runLock.unlock();
} catch (InterruptedException e) {
// Ignored
}
private boolean stopRunningLoop() {
thread.stop(true);
return thread.awaitTermination(1000);
}
private void queueConnectionLoop(long delay) {
thread.execute(() -> startConnectionLoop(delay));
}
private void startConnectionLoop(long delay) {
if (sleepLock.tryLock()) {
try {
if (!Delay.sleepMilli(delay))
private void primaryConnectionLoop() {
try (HolocoreSocket holocore = new HolocoreSocket(data.getAddress().getAddress(), data.getAddress().getPort())) {
this.holocore = holocore;
Log.t("Attempting to connect to server at %s", holocore.getRemoteAddress());
if (!holocore.connect(CONNECT_TIMEOUT)) {
Log.w("Failed to connect to server!");
return;
}
Log.i("Successfully connected to server at %s", holocore.getRemoteAddress());
intentChain.broadcastAfter(getIntentManager(), new ServerConnectedIntent());
while (holocore.isConnected()) {
RawPacket inbound = holocore.receive();
if (inbound == null) {
Log.w("Server closed connection!");
return;
} finally {
sleepLock.unlock();
}
} else if (delay > 0) {
return;
}
if (runLock.tryLock()) {
try {
activeThread.set(Thread.currentThread());
holocore = new HolocoreSocket(data.getAddress().getAddress(), data.getAddress().getPort());
running.set(true);
if (attemptConnection()) {
while (holocore.isConnected()) {
if (!connectionLoop())
break;
}
}
} catch (Throwable t) {
Log.w(t);
Log.i("Disconnected from server. Sleeping 3 seconds");
holocore.disconnect(ServerConnectionChangedReason.UNKNOWN);
queueConnectionLoop(3000);
} finally {
cleanupConnection();
activeThread.set(null);
runLock.unlock();
intentChain.broadcastAfter(getIntentManager(), new DataPacketOutboundIntent(interceptor.interceptServer(inbound.getData())));
}
} catch (Throwable t) {
Log.w("Caught unknown exception in server connection! %s: %s", t.getClass().getName(), t.getMessage());
Log.w(t);
} finally {
Log.i("Disconnected from server.");
intentChain.broadcastAfter(getIntentManager(), new ServerDisconnectedIntent());
this.holocore = null;
}
}
private boolean attemptConnection() {
Log.t("Attempting to connect to server at %s", holocore.getRemoteAddress());
if (!holocore.connect(5000)) {
Log.t("Failed to connect to server. Sleeping 3 seconds");
queueConnectionLoop(3000);
return false;
}
intentChain.broadcastAfter(getIntentManager(), new ServerConnectedIntent());
Log.i("Successfully connected to server at %s", holocore.getRemoteAddress());
return true;
}
private boolean connectionLoop() {
if (!running.get()) {
holocore.disconnect(ServerConnectionChangedReason.CLIENT_DISCONNECT);
return false;
}
RawPacket inbound = holocore.receive();
if (inbound == null) {
holocore.disconnect(ServerConnectionChangedReason.SOCKET_CLOSED);
return false;
}
intentChain.broadcastAfter(getIntentManager(), new DataPacketOutboundIntent(interceptor.interceptServer(inbound.getData())));
return true;
}
private void cleanupConnection() {
intentChain.broadcastAfter(getIntentManager(), new ServerDisconnectedIntent());
Log.i("Disconnected from server");
holocore.terminate();
holocore = null;
running.set(false);
}
}