From 668025fac8b4ea7361e26e6407c43268f19eb12a Mon Sep 17 00:00:00 2001 From: Josh Larson Date: Thu, 30 Aug 2018 15:57:18 -0500 Subject: [PATCH] Various bugfixes and improvements --- client-holocore | 2 +- pswgcommon | 2 +- .../com/projectswg/forwarder/Forwarder.java | 25 ++- .../server/RequestServerConnectionIntent.java | 31 +++ .../networking/data/FragmentedProcessor.java | 31 ++- .../resources/networking/data/Packager.java | 30 ++- .../networking/data/ProtocolStack.java | 168 +++++++++------- .../networking/data/SequencedOutbound.java | 30 ++- .../networking/packets/DataChannel.java | 32 ++- .../networking/packets/Fragmented.java | 139 ++++++------- .../networking/packets/SequencedPacket.java | 29 +++ .../client/ClientInboundDataService.java | 44 +++-- .../client/ClientOutboundDataService.java | 185 +++++++++--------- .../services/client/ClientServerService.java | 39 +++- .../server/ServerConnectionService.java | 153 ++++----------- 15 files changed, 499 insertions(+), 441 deletions(-) create mode 100644 src/main/java/com/projectswg/forwarder/intents/server/RequestServerConnectionIntent.java diff --git a/client-holocore b/client-holocore index 450ce04..5ca8bf2 160000 --- a/client-holocore +++ b/client-holocore @@ -1 +1 @@ -Subproject commit 450ce04c13ac4f80bb5061bb38a2abfe0939ed93 +Subproject commit 5ca8bf29ac55824735723713729e713d1f36c50d diff --git a/pswgcommon b/pswgcommon index e2532b3..81bd089 160000 --- a/pswgcommon +++ b/pswgcommon @@ -1 +1 @@ -Subproject commit e2532b34eda75626a06ace134211bd19c3ccbde5 +Subproject commit 81bd0897d977eaa51f2b8cc7b1c9a2c1986dedca diff --git a/src/main/java/com/projectswg/forwarder/Forwarder.java b/src/main/java/com/projectswg/forwarder/Forwarder.java index dc3ea8f..4d3b703 100644 --- a/src/main/java/com/projectswg/forwarder/Forwarder.java +++ b/src/main/java/com/projectswg/forwarder/Forwarder.java @@ -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 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() { diff --git a/src/main/java/com/projectswg/forwarder/intents/server/RequestServerConnectionIntent.java b/src/main/java/com/projectswg/forwarder/intents/server/RequestServerConnectionIntent.java new file mode 100644 index 0000000..7480873 --- /dev/null +++ b/src/main/java/com/projectswg/forwarder/intents/server/RequestServerConnectionIntent.java @@ -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 . * + * * + ***********************************************************************************/ + +package com.projectswg.forwarder.intents.server; + +import me.joshlarson.jlcommon.control.Intent; + +public class RequestServerConnectionIntent extends Intent { + + public RequestServerConnectionIntent() { + + } + +} diff --git a/src/main/java/com/projectswg/forwarder/resources/networking/data/FragmentedProcessor.java b/src/main/java/com/projectswg/forwarder/resources/networking/data/FragmentedProcessor.java index 2a92ebf..c6839e0 100644 --- a/src/main/java/com/projectswg/forwarder/resources/networking/data/FragmentedProcessor.java +++ b/src/main/java/com/projectswg/forwarder/resources/networking/data/FragmentedProcessor.java @@ -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; } diff --git a/src/main/java/com/projectswg/forwarder/resources/networking/data/Packager.java b/src/main/java/com/projectswg/forwarder/resources/networking/data/Packager.java index 0fd2943..617dc6e 100644 --- a/src/main/java/com/projectswg/forwarder/resources/networking/data/Packager.java +++ b/src/main/java/com/projectswg/forwarder/resources/networking/data/Packager.java @@ -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 dataChannel; private final Queue outboundRaw; - private final Queue outboundPackaged; - private final ProtocolStack stack; + private final ConnectionStream outboundPackaged; - public Packager(Queue outboundRaw, Queue outboundPackaged, ProtocolStack stack) { + public Packager(Queue outboundRaw, ConnectionStream 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); } diff --git a/src/main/java/com/projectswg/forwarder/resources/networking/data/ProtocolStack.java b/src/main/java/com/projectswg/forwarder/resources/networking/data/ProtocolStack.java index 0254d30..59a7dee 100644 --- a/src/main/java/com/projectswg/forwarder/resources/networking/data/ProtocolStack.java +++ b/src/main/java/com/projectswg/forwarder/resources/networking/data/ProtocolStack.java @@ -13,37 +13,29 @@ import java.util.function.BiConsumer; public class ProtocolStack { - private final PriorityQueue sequenced; private final FragmentedProcessor fragmentedProcessor; private final InetSocketAddress source; private final BiConsumer sender; private final ClientServer server; private final Queue outboundRaw; - private final Queue outboundPackaged; + private final ConnectionStream inbound; + private final ConnectionStream 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 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 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 { + + private final PriorityQueue sequenced; + private final PriorityQueue 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 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 + } + } diff --git a/src/main/java/com/projectswg/forwarder/resources/networking/data/SequencedOutbound.java b/src/main/java/com/projectswg/forwarder/resources/networking/data/SequencedOutbound.java index 6f01690..2a5d067 100644 --- a/src/main/java/com/projectswg/forwarder/resources/networking/data/SequencedOutbound.java +++ b/src/main/java/com/projectswg/forwarder/resources/networking/data/SequencedOutbound.java @@ -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; } diff --git a/src/main/java/com/projectswg/forwarder/resources/networking/packets/DataChannel.java b/src/main/java/com/projectswg/forwarder/resources/networking/packets/DataChannel.java index a5bd9ad..c0c1d68 100644 --- a/src/main/java/com/projectswg/forwarder/resources/networking/packets/DataChannel.java +++ b/src/main/java/com/projectswg/forwarder/resources/networking/packets/DataChannel.java @@ -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 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 diff --git a/src/main/java/com/projectswg/forwarder/resources/networking/packets/Fragmented.java b/src/main/java/com/projectswg/forwarder/resources/networking/packets/Fragmented.java index c7538c1..f6b3b8b 100644 --- a/src/main/java/com/projectswg/forwarder/resources/networking/packets/Fragmented.java +++ b/src/main/java/com/projectswg/forwarder/resources/networking/packets/Fragmented.java @@ -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); - } - } diff --git a/src/main/java/com/projectswg/forwarder/resources/networking/packets/SequencedPacket.java b/src/main/java/com/projectswg/forwarder/resources/networking/packets/SequencedPacket.java index 3c432d7..688b782 100644 --- a/src/main/java/com/projectswg/forwarder/resources/networking/packets/SequencedPacket.java +++ b/src/main/java/com/projectswg/forwarder/resources/networking/packets/SequencedPacket.java @@ -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 { + 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; + } + } diff --git a/src/main/java/com/projectswg/forwarder/services/client/ClientInboundDataService.java b/src/main/java/com/projectswg/forwarder/services/client/ClientInboundDataService.java index eb7ebc2..9a26fdc 100644 --- a/src/main/java/com/projectswg/forwarder/services/client/ClientInboundDataService.java +++ b/src/main/java/com/projectswg/forwarder/services/client/ClientInboundDataService.java @@ -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)); } diff --git a/src/main/java/com/projectswg/forwarder/services/client/ClientOutboundDataService.java b/src/main/java/com/projectswg/forwarder/services/client/ClientOutboundDataService.java index e0f487f..1272505 100644 --- a/src/main/java/com/projectswg/forwarder/services/client/ClientOutboundDataService.java +++ b/src/main/java/com/projectswg/forwarder/services/client/ClientOutboundDataService.java @@ -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 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"); } } diff --git a/src/main/java/com/projectswg/forwarder/services/client/ClientServerService.java b/src/main/java/com/projectswg/forwarder/services/client/ClientServerService.java index 55c64a7..2f2f0ea 100644 --- a/src/main/java/com/projectswg/forwarder/services/client/ClientServerService.java +++ b/src/main/java/com/projectswg/forwarder/services/client/ClientServerService.java @@ -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 stacks; + private final AtomicBoolean serverConnection; + private final AtomicReference 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()); diff --git a/src/main/java/com/projectswg/forwarder/services/server/ServerConnectionService.java b/src/main/java/com/projectswg/forwarder/services/server/ServerConnectionService.java index 43980fd..38ec17c 100644 --- a/src/main/java/com/projectswg/forwarder/services/server/ServerConnectionService.java +++ b/src/main/java/com/projectswg/forwarder/services/server/ServerConnectionService.java @@ -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 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); - } - }