Optimized Packet Writing (separate Thread for I/O operations)

This commit is contained in:
Marius Metzger
2015-01-04 21:51:56 +01:00
parent 7adc9b86c5
commit 5f4d12cb81
9 changed files with 152 additions and 422 deletions

View File

@@ -54,7 +54,7 @@ public class GuiReplayOverlay extends Gui {
private int ppButtonX = 10;
private int ppButtonY = 10;
private int r_ppButtonX = 10;
private int r_ppButtonY = realTimelineY+1;
@@ -191,8 +191,10 @@ public class GuiReplayOverlay extends Gui {
}
} else {
speedSlider.mouseReleased(mouseX, mouseY);
mouseDown = false;
try {
speedSlider.mouseReleased(mouseX, mouseY);
mouseDown = false;
} catch(Exception e) {}
}
@@ -605,10 +607,10 @@ public class GuiReplayOverlay extends Gui {
this.drawModalRectWithCustomSizedTexture((int)real_x, ry, dx, dy, 5, 5, 64, 64);
}
}
//Draw Play/Pause Button
//Play/Pause button
int dx = 0;
int dy = 0;
@@ -633,7 +635,7 @@ public class GuiReplayOverlay extends Gui {
GlStateManager.resetColor();
this.drawModalRectWithCustomSizedTexture(r_ppButtonX, r_ppButtonY, dx, dy, 20, 20, 64, 64);
//Handling the click on the Replay starter
if(hover && Mouse.isButtonDown(0) && isClick()) {
if(ReplayHandler.isReplaying()) {

View File

@@ -1,14 +1,14 @@
package eu.crushedpixel.replaymod.unused;
package eu.crushedpixel.replaymod.holders;
import io.netty.buffer.ByteBuf;
import net.minecraft.network.Packet;
public class DataReciever {
public class PacketData {
private byte[] array;
private int timestamp;
public DataReciever(byte[] array, int timestamp) {
public PacketData(byte[] array, int timestamp) {
this.array = array;
this.timestamp = timestamp;
}

View File

@@ -35,15 +35,15 @@ public class ConnectionEventHandler {
private File currentFile;
private String fileName;
private static PacketListener packetListener = null;
private static boolean isRecording = false;
public static boolean isRecording() {
return isRecording;
}
public static void insertPacket(Packet packet) {
if(!isRecording || packetListener == null) {
System.out.println("Invalid attempt to insert Packet!");
@@ -55,7 +55,7 @@ public class ConnectionEventHandler {
e.printStackTrace();
}
}
@SubscribeEvent
public void onConnectedToServerEvent(ClientConnectedToServerEvent event) {
System.out.println("Connected to server");
@@ -101,17 +101,10 @@ public class ConnectionEventHandler {
PacketListener insert = null;
if(event.isLocal) {
pipeline.addBefore(packetHandlerKey, "replay_recorder", insert = new PacketListener(currentFile, fileName, worldName, System.currentTimeMillis(), maxFileSize));
ChatMessageRequests.addChatMessage("Recording started!", ChatMessageType.INFORMATION);
isRecording = true;
} else {
//pipeline.addBefore(decoderKey, "replay_recorder", insert = new DataListener(currentFile, fileName, worldName, System.currentTimeMillis(), maxFileSize));
pipeline.addBefore(packetHandlerKey, "replay_recorder", insert = new PacketListener(currentFile, fileName, worldName, System.currentTimeMillis(), maxFileSize));
ChatMessageRequests.addChatMessage("Recording started!", ChatMessageType.INFORMATION);
isRecording = true;
}
pipeline.addBefore(packetHandlerKey, "replay_recorder", insert = new PacketListener
(currentFile, fileName, worldName, System.currentTimeMillis(), maxFileSize, event.isLocal));
ChatMessageRequests.addChatMessage("Recording started!", ChatMessageType.INFORMATION);
isRecording = true;
final PacketListener listener = insert;
@@ -139,7 +132,7 @@ public class ConnectionEventHandler {
}
}).start();
}
packetListener = listener;
} catch(Exception e) {

View File

@@ -11,6 +11,9 @@ import java.io.FileInputStream;
import java.io.FileNotFoundException;
import java.io.FileOutputStream;
import java.io.PrintWriter;
import java.util.HashMap;
import java.util.Map.Entry;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.zip.ZipEntry;
import java.util.zip.ZipOutputStream;
@@ -20,23 +23,24 @@ import com.google.gson.Gson;
import eu.crushedpixel.replaymod.chat.ChatMessageRequests;
import eu.crushedpixel.replaymod.chat.ChatMessageRequests.ChatMessageType;
import eu.crushedpixel.replaymod.holders.PacketData;
public class DataListener extends ChannelInboundHandlerAdapter {
public abstract class DataListener extends ChannelInboundHandlerAdapter {
protected File file;
protected long startTime;
protected Long startTime = null;
protected long maxSize;
protected long totalBytes = 0;
protected DataOutputStream out;
protected String name;
protected String worldName;
private boolean singleplayer;
protected long lastSentPacket = 0;
protected boolean alive = true;
protected boolean isPacketListener = false;
private boolean writeJoinPacket = false;
protected DataWriter dataWriter;
private Gson gson = new Gson();
@@ -44,141 +48,136 @@ public class DataListener extends ChannelInboundHandlerAdapter {
this.worldName = worldName;
System.out.println(worldName);
}
public void insertPacket(Packet p) {
}
public DataListener(File file, String name, String worldName, long startTime, int maxSize) throws FileNotFoundException {
public DataListener(File file, String name, String worldName, long startTime, int maxSize, boolean singleplayer) throws FileNotFoundException {
this.file = file;
this.startTime = startTime;
this.maxSize = maxSize*1024*1024;
this.name = name;
this.worldName = worldName;
this.singleplayer = singleplayer;
System.out.println(worldName);
FileOutputStream fos = new FileOutputStream(file);
BufferedOutputStream bos = new BufferedOutputStream(fos);
out = new DataOutputStream(bos);
writeJoinPacket = true;
DataOutputStream out = new DataOutputStream(bos);
dataWriter = new DataWriter(out);
}
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
if(!alive || isPacketListener) {
lastSentPacket = System.currentTimeMillis();
super.channelRead(ctx, msg);
return;
}
/*
if(writeJoinPacket) {
writeJoinPacket = false;
S01PacketJoinGame p = new S01PacketJoinGame();
int timestamp = 0;
PacketSerializer ps = new PacketSerializer(EnumPacketDirection.CLIENTBOUND);
ByteBuf bb = Unpooled.buffer();
PacketBuffer pb = new PacketBuffer(bb);
pb.writeInt(new Random().nextInt(1000));
pb.writeByte(3);
pb.writeByte(-1);
pb.writeByte(0);
pb.writeByte(100);
pb.writeString("default");
pb.writeBoolean(false);
p.readPacketData(pb);
ByteBuf b2 = Unpooled.buffer();
ps.encode(ctx, p, b2);
b2.readerIndex(0);
byte[] array = new byte[b2.readableBytes()];
b2.readBytes(array);
b2.readerIndex(0);
out.writeInt(timestamp);
out.writeInt(array.length);
out.write(array);
out.flush();
}
*/
int timestamp = (int)(System.currentTimeMillis() - startTime);
ByteBuf byteBuf = (ByteBuf)msg;
byteBuf.readerIndex(0);
byte[] array = new byte[byteBuf.readableBytes()];
byteBuf.readBytes(array);
totalBytes += (array.length + (2*4)); //two Integer values and the Packet size
if(totalBytes >= maxSize && maxSize > 0) {
ChatMessageRequests.addChatMessage("Maximum file size exceeded", ChatMessageType.WARNING);
ChatMessageRequests.addChatMessage("The Recording has been stopped", ChatMessageType.WARNING);
alive = false;
}
byteBuf.readerIndex(0);
out.writeInt(timestamp);
out.writeInt(array.length);
out.write(array);
out.flush();
super.channelRead(ctx, msg);
}
@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception {
out.flush();
out.close();
dataWriter.requestFinish();
}
public class DataWriter {
byte[] buffer = new byte[1024];
private boolean active = true;
try {
ReplayMetaData metaData = new ReplayMetaData(isPacketListener, worldName, (int) (lastSentPacket - startTime), startTime);
String json = gson.toJson(metaData);
private ConcurrentLinkedQueue<PacketData> queue = new ConcurrentLinkedQueue<PacketData>();
File folder = new File("./replay_recordings/");
folder.mkdirs();
File archive = new File(folder, name+ConnectionEventHandler.ZIP_FILE_EXTENSION);
archive.createNewFile();
FileOutputStream fos = new FileOutputStream(archive);
ZipOutputStream zos = new ZipOutputStream(fos);
zos.putNextEntry(new ZipEntry("metaData.json"));
PrintWriter pw = new PrintWriter(zos);
pw.write(json);
pw.flush();
zos.closeEntry();
zos.putNextEntry(new ZipEntry("recording"+ConnectionEventHandler.TEMP_FILE_EXTENSION));
FileInputStream fis = new FileInputStream(file);
int len;
while((len = fis.read(buffer)) > 0) {
zos.write(buffer, 0, len);
}
fis.close();
zos.closeEntry();
zos.close();
file.delete();
} catch(Exception e) {
e.printStackTrace();
public void writeData(PacketData data) {
queue.add(data);
}
Thread outputThread = new Thread(new Runnable() {
@Override
public void run() {
HashMap<Class, Integer> counts = new HashMap<Class, Integer>();
while(active) {
PacketData dataReciever = queue.poll();
if(dataReciever != null) {
//write the ByteBuf to the given OutputStream
byte[] array = dataReciever.getByteArray();
if(array != null) {
try {
stream.writeInt(dataReciever.getTimestamp()); //Timestamp
stream.writeInt(array.length); //Lenght
stream.write(array); //Content
stream.flush();
} catch(Exception e) {
e.printStackTrace();
}
}
} else {
try {
//let the Thread sleep for 1/4 second and queue up new Packets
Thread.sleep(250L);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
try {
stream.flush();
stream.close();
} catch (Exception e) {
e.printStackTrace();
}
for(Entry<Class, Integer> entries : counts.entrySet()) {
System.out.println(entries.getKey()+ "| "+entries.getValue());
}
}
});
private DataOutputStream stream;
public DataWriter(DataOutputStream stream) {
this.stream = stream;
outputThread.start();
}
public void requestFinish() {
active = false;
byte[] buffer = new byte[1024];
try {
ReplayMetaData metaData = new ReplayMetaData(singleplayer, worldName, (int) lastSentPacket, startTime);
String json = gson.toJson(metaData);
File folder = new File("./replay_recordings/");
folder.mkdirs();
File archive = new File(folder, name+ConnectionEventHandler.ZIP_FILE_EXTENSION);
archive.createNewFile();
FileOutputStream fos = new FileOutputStream(archive);
ZipOutputStream zos = new ZipOutputStream(fos);
zos.putNextEntry(new ZipEntry("metaData.json"));
PrintWriter pw = new PrintWriter(zos);
pw.write(json);
pw.flush();
zos.closeEntry();
zos.putNextEntry(new ZipEntry("recording"+ConnectionEventHandler.TEMP_FILE_EXTENSION));
FileInputStream fis = new FileInputStream(file);
int len;
while((len = fis.read(buffer)) > 0) {
zos.write(buffer, 0, len);
}
fis.close();
zos.closeEntry();
zos.close();
file.delete();
} catch(Exception e) {
e.printStackTrace();
}
}
}
}

View File

@@ -15,17 +15,19 @@ import net.minecraft.network.play.server.S0CPacketSpawnPlayer;
import net.minecraft.network.play.server.S0FPacketSpawnMob;
import eu.crushedpixel.replaymod.chat.ChatMessageRequests;
import eu.crushedpixel.replaymod.chat.ChatMessageRequests.ChatMessageType;
import eu.crushedpixel.replaymod.holders.PacketData;
import eu.crushedpixel.replaymod.reflection.MCPNames;
public class PacketListener extends DataListener {
public PacketListener(File file, String name, String worldName, long startTime, int maxSize) throws FileNotFoundException {
super(file, name, worldName, startTime, maxSize);
isPacketListener = true;
public PacketListener(File file, String name, String worldName, long startTime, int maxSize, boolean singleplayer) throws FileNotFoundException {
super(file, name, worldName, startTime, maxSize, singleplayer);
}
private ChannelHandlerContext context = null;
private static final PacketSerializer packetSerializer = new PacketSerializer(EnumPacketDirection.CLIENTBOUND);
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
if(ctx == null) {
@@ -44,12 +46,12 @@ public class PacketListener extends DataListener {
try {
Packet packet = (Packet)msg;
if(startTime == null) startTime = System.currentTimeMillis();
int timestamp = (int)(System.currentTimeMillis() - startTime);
//Converts the packet back to a ByteBuffer for correct saving
PacketSerializer ps = new PacketSerializer(EnumPacketDirection.CLIENTBOUND);
ByteBuf bb = Unpooled.buffer();
if(packet instanceof S0FPacketSpawnMob) {
Field field_149043_l = S0FPacketSpawnMob.class.getDeclaredField(MCPNames.field("field_149043_l"));
@@ -71,13 +73,12 @@ public class PacketListener extends DataListener {
}
}
ps.encode(ctx, packet, bb);
packetSerializer.encode(ctx, packet, bb);
bb.readerIndex(0);
byte[] array = new byte[bb.readableBytes()];
bb.readBytes(array);
totalBytes += (array.length + (2*4)); //two Integer values and the Packet size
if(totalBytes >= maxSize && maxSize > 0) {
ChatMessageRequests.addChatMessage("Maximum file size exceeded", ChatMessageType.WARNING);
@@ -87,11 +88,9 @@ public class PacketListener extends DataListener {
bb.readerIndex(0);
out.writeInt(timestamp); //Timestamp
out.writeInt(array.length); //Lenght
out.write(array);
out.flush();
dataWriter.writeData(new PacketData(array, timestamp));
lastSentPacket = timestamp;
} catch(Exception e) {
e.printStackTrace();
}

View File

@@ -1,88 +0,0 @@
package eu.crushedpixel.replaymod.unused;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import java.io.DataOutputStream;
import java.util.HashMap;
import java.util.Map.Entry;
import java.util.concurrent.ConcurrentLinkedQueue;
import net.minecraft.network.PacketBuffer;
public class DataWriter {
private boolean active = true;
private ConcurrentLinkedQueue<DataReciever> queue = new ConcurrentLinkedQueue<DataReciever>();
Thread outputThread = new Thread(new Runnable() {
@Override
public void run() {
HashMap<Class, Integer> counts = new HashMap<Class, Integer>();
while(active) {
DataReciever dataReciever = queue.poll();
if(dataReciever != null) {
//write the ByteBuf to the given OutputStream
byte[] array = dataReciever.getByteArray();
if(array != null) {
try {
stream.writeInt(dataReciever.getTimestamp()); //Timestamp
stream.writeInt(array.length); //Lenght
stream.write(array); //Content
stream.flush();
} catch(Exception e) {
e.printStackTrace();
}
}
} else {
try {
//let the Thread sleep for 1/2 second and queue up new Packets
Thread.sleep(500L);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
try {
stream.close();
} catch (Exception e) {
e.printStackTrace();
}
for(Entry<Class, Integer> entries : counts.entrySet()) {
System.out.println(entries.getKey()+ "| "+entries.getValue());
}
}
});
public ConcurrentLinkedQueue<DataReciever> getQueue() {
return queue;
}
public void setQueue(ConcurrentLinkedQueue<DataReciever> queue) {
this.queue = queue;
}
private DataOutputStream stream;
public DataWriter(DataOutputStream stream) {
this.stream = stream;
outputThread.start();
}
public void requestFinish() {
active = false;
}
}

View File

@@ -1,37 +0,0 @@
package eu.crushedpixel.replaymod.unused;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import net.minecraft.network.EnumConnectionState;
import net.minecraft.network.EnumPacketDirection;
import net.minecraft.network.NetworkManager;
import net.minecraft.network.Packet;
import net.minecraft.network.PacketBuffer;
public class PacketListener extends ChannelInboundHandlerAdapter {
private PacketWriter packetWriter;
private long startTime;
public PacketListener(PacketWriter packetWriter, long startTime) {
this.packetWriter = packetWriter;
this.startTime = startTime;
}
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
super.channelRead(ctx, msg);
if(msg instanceof Packet) {
Packet packet = (Packet)msg;
int packetID = ((EnumConnectionState)ctx.channel()
.attr(NetworkManager.attrKeyConnectionState).get()).getPacketId(EnumPacketDirection.CLIENTBOUND, packet);
int timestamp = (int) (System.currentTimeMillis() - startTime);
PacketReciever reciever = new PacketReciever(packet, (short)packetID, timestamp);
packetWriter.getQueue().add(reciever);
}
}
}

View File

@@ -1,37 +0,0 @@
package eu.crushedpixel.replaymod.unused;
import net.minecraft.network.Packet;
public class PacketReciever {
private Packet packet;
private short packetID;
private int timestamp;
public PacketReciever(Packet packet, short packetID, int timestamp) {
this.packet = packet;
this.packetID = packetID;
this.timestamp = timestamp;
}
public Packet getPacket() {
return packet;
}
public void setPacket(Packet packet) {
this.packet = packet;
}
public short getPacketID() {
return packetID;
}
public void setPacketID(short packetID) {
this.packetID = packetID;
}
public int getTimestamp() {
return timestamp;
}
public void setTimestamp(int timestamp) {
this.timestamp = timestamp;
}
}

View File

@@ -1,101 +0,0 @@
package eu.crushedpixel.replaymod.unused;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import java.io.DataOutputStream;
import java.util.HashMap;
import java.util.Map.Entry;
import java.util.concurrent.ConcurrentLinkedQueue;
import net.minecraft.network.PacketBuffer;
public class PacketWriter {
private boolean active = true;
private ConcurrentLinkedQueue<PacketReciever> queue = new ConcurrentLinkedQueue<PacketReciever>();
Thread outputThread = new Thread(new Runnable() {
@Override
public void run() {
HashMap<Class, Integer> counts = new HashMap<Class, Integer>();
while(active) {
PacketReciever packetReciever = queue.poll();
if(packetReciever != null) {
ByteBuf bb = Unpooled.buffer();
PacketBuffer packetBuffer = new PacketBuffer(bb);
//write the Packet to the given OutputStream
if(packetReciever.getPacket() != null) {
try {
packetReciever.getPacket().writePacketData(packetBuffer);
bb.readerIndex(0);
byte[] array = new byte[bb.readableBytes()];
bb.readBytes(array);
stream.writeInt(packetReciever.getTimestamp()); //Timestamp
stream.writeInt(array.length); //Lenght
if(counts.containsKey(packetReciever.getPacket().getClass())) {
counts.put(packetReciever.getPacket().getClass(), counts.get(packetReciever.getPacket().getClass())+array.length);
} else {
counts.put(packetReciever.getPacket().getClass(), array.length);
}
stream.writeShort(packetReciever.getPacketID()); //Packet ID
stream.write(array); //Content
stream.flush();
} catch(Exception e) {
System.out.println(packetReciever.getPacket().getClass());
e.printStackTrace();
}
}
} else {
try {
//let the Thread sleep for 1/2 second and queue up new Packets
Thread.sleep(500L);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
try {
stream.close();
} catch (Exception e) {
e.printStackTrace();
}
for(Entry<Class, Integer> entries : counts.entrySet()) {
System.out.println(entries.getKey()+ "| "+entries.getValue());
}
}
});
public ConcurrentLinkedQueue<PacketReciever> getQueue() {
return queue;
}
public void setQueue(ConcurrentLinkedQueue<PacketReciever> queue) {
this.queue = queue;
}
private DataOutputStream stream;
public PacketWriter(DataOutputStream stream) {
this.stream = stream;
outputThread.start();
}
public void requestFinish() {
active = false;
}
}