diff --git a/src/main/java/com/replaymod/replay/QuickReplaySender.java b/src/main/java/com/replaymod/replay/QuickReplaySender.java index 31b376cf..51a12bf0 100644 --- a/src/main/java/com/replaymod/replay/QuickReplaySender.java +++ b/src/main/java/com/replaymod/replay/QuickReplaySender.java @@ -32,6 +32,8 @@ import com.github.steveice10.packetlib.io.NetInput; import com.github.steveice10.packetlib.io.NetOutput; import com.github.steveice10.packetlib.io.stream.StreamNetInput; import com.github.steveice10.packetlib.io.stream.StreamNetOutput; +import com.github.steveice10.packetlib.tcp.io.ByteBufNetInput; +import com.github.steveice10.packetlib.tcp.io.ByteBufNetOutput; import com.google.common.base.Optional; import com.google.common.collect.ListMultimap; import com.google.common.collect.Multimaps; @@ -53,6 +55,7 @@ import io.netty.buffer.Unpooled; import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerAdapter; import io.netty.channel.ChannelHandlerContext; +import lombok.SneakyThrows; import net.minecraft.client.Minecraft; import net.minecraft.network.EnumConnectionState; import net.minecraft.network.EnumPacketDirection; @@ -95,6 +98,7 @@ import static com.replaymod.replay.ReplayModReplay.LOGGER; @ChannelHandler.Sharable public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySender { private static final String CACHE_ENTRY = "quickModeCache.bin"; + private static final String CACHE_INDEX_ENTRY = "quickModeCacheIndex.bin"; private static final int CACHE_VERSION = 0; private final Minecraft mc = Minecraft.getMinecraft(); @@ -115,6 +119,8 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe private ListenableFuture initPromise; + private com.github.steveice10.netty.buffer.ByteBuf buf; + private NetInput bufInput; private TreeMap> thingSpawnsT = new TreeMap<>(); private ListMultimap thingSpawns = Multimaps.newListMultimap(thingSpawnsT, ArrayList::new); private TreeMap> thingDespawnsT = new TreeMap<>(); @@ -207,7 +213,7 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe public void setReplaySpeed(double factor) { if (factor != 0) { if (paused() && asyncMode) { - lastAsyncUpdateTime = System.currentTimeMillis(); // TODO test this + lastAsyncUpdateTime = System.currentTimeMillis(); } this.replaySpeed = factor; } @@ -263,49 +269,63 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } private boolean tryLoadFromCache(Consumer progress) throws IOException { - Optional cacheOpt = replayFile.get(CACHE_ENTRY); - if (!cacheOpt.isPresent()) return false; - boolean success = loadFromCache(cacheOpt.get(), progress); - if (!success) { - thingSpawnsT.clear(); - thingDespawnsT.clear(); - worldTimes.clear(); - thunderStrengths.clear(); - } - return success; - } + boolean success = false; - private boolean loadFromCache(InputStream rawIn, Consumer progress) throws IOException { - long sysTimeStart = System.currentTimeMillis(); - NetInput in = new StreamNetInput(rawIn); - int version = in.readVarInt(); - if (version > CACHE_VERSION) return false; // More recent than this version, no chance we can read that - - int size = in.readVarInt(); - int spawnTime = 0; - for (int i = 0; i < size; i++) { - progress.accept((double) i / size); - spawnTime += in.readVarInt(); - for (int j = in.readVarInt(); j > 0; j--) { - int despawnTime = in.readVarInt(); - BakedTrackedThing trackedThing; - switch (in.readVarInt()) { - case 0: trackedThing = new BakedEntity(version, in); break; - case 1: trackedThing = new BakedChunk(version, in); break; - case 2: trackedThing = new BakedWeather(version, in); break; - default: return false; - } - trackedThing.spawnTime = spawnTime; - trackedThing.despawnTime = despawnTime; - thingSpawns.put(spawnTime, trackedThing); - thingDespawns.put(despawnTime, trackedThing); + Optional cacheIndexOpt = replayFile.get(CACHE_INDEX_ENTRY); + if (!cacheIndexOpt.isPresent()) return false; + try (InputStream indexIn = cacheIndexOpt.get()) { + Optional cacheOpt = replayFile.get(CACHE_ENTRY); + if (!cacheOpt.isPresent()) return false; + try (InputStream cacheIn = cacheOpt.get()) { + success = loadFromCache(cacheIn, indexIn, progress); + } + } finally { + if (!success) { + buf = null; + bufInput = null; + thingSpawnsT.clear(); + thingDespawnsT.clear(); + worldTimes.clear(); + thunderStrengths.clear(); } } - // These should go quick, they probably won't need to update progress - progress.accept(1.0); - readFromCache(version, in, worldTimes); - readFromCache(version, in, thunderStrengths); + return success; + } + + private boolean loadFromCache(InputStream cacheIn, InputStream rawIndexIn, Consumer progress) throws IOException { + long sysTimeStart = System.currentTimeMillis(); + + NetInput in = new StreamNetInput(rawIndexIn); + if (in.readVarInt() != CACHE_VERSION) return false; // Incompatible cache version + if (new StreamNetInput(cacheIn).readVarInt() != CACHE_VERSION) return false; // Incompatible cache version + + things: while (true) { + BakedTrackedThing trackedThing; + switch (in.readVarInt()) { + case 0: break things; + case 1: trackedThing = new BakedEntity(in); break; + case 2: trackedThing = new BakedChunk(in); break; + case 3: trackedThing = new BakedWeather(in); break; + default: return false; + } + thingSpawns.put(trackedThing.spawnTime, trackedThing); + thingDespawns.put(trackedThing.despawnTime, trackedThing); + } + + readFromCache(in, worldTimes); + readFromCache(in, thunderStrengths); + int size = in.readVarInt(); + + buf = com.github.steveice10.netty.buffer.Unpooled.buffer(size); + int read = 0; + while (true) { + int len = buf.writeBytes(cacheIn, Math.min(size - read, 4096)); + if (len <= 0) break; + read += len; + progress.accept((double) read / size); + } + bufInput = new ByteBufNetInput(buf); LOGGER.info("Loaded quick replay from cache in " + (System.currentTimeMillis() - sysTimeStart) + "ms"); return true; @@ -327,8 +347,6 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe studio.setParsing(ServerUpdateTimePacket.class, true); studio.setParsing(ServerNotifyClientPacket.class, true); - TreeMap> thingSpawnsT = new TreeMap<>(); - ListMultimap thingSpawns = Multimaps.newListMultimap(thingSpawnsT, ArrayList::new); TreeMap worldTimes = new TreeMap<>(); TreeMap thunderStrengths = new TreeMap<>(); Map playerListEntries = new HashMap<>(); @@ -336,33 +354,33 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe Map activeChunks = new HashMap<>(); Weather activeWeather = null; - // 90% of estimated time are for analysis, 10% for writing cache - // This heuristic is based on simply running the function on some 3h replay. - // Numbers may be significantly different for smaller replays where GC doesn't kick in. But the progress bar - // will hopefully be barely visible in that case anyway. - // May need to be updated when performance of MCProtocolLib or this code is improved. - // Same for the split between analyseReplay and loadFromCache declared in #initialize(...) - final double progressSplit = 0.9; - double sysTimeStart = System.currentTimeMillis(); double duration; - try (ReplayInputStream in = replayFile.getPacketData(studio)) { + try (ReplayInputStream in = replayFile.getPacketData(studio); + OutputStream cacheOut = replayFile.write(CACHE_ENTRY); + OutputStream cacheIndexOut = replayFile.write(CACHE_INDEX_ENTRY)) { + NetOutput out = new StreamNetOutput(cacheOut); + out.writeVarInt(CACHE_VERSION); + NetOutput indexOut = new StreamNetOutput(cacheIndexOut); + indexOut.writeVarInt(CACHE_VERSION); + + int index = 0; + int time = 0; duration = replayFile.getMetaData().getDuration(); PacketData packetData; while ((packetData = in.readPacket()) != null) { com.github.steveice10.packetlib.packet.Packet packet = packetData.getPacket(); - int time = (int) packetData.getTime(); - progress.accept(time / duration * progressSplit); + time = (int) packetData.getTime(); + progress.accept(time / duration); Integer entityId = PacketUtils.getEntityId(packet); if (packet instanceof ServerSpawnMobPacket || packet instanceof ServerSpawnObjectPacket || packet instanceof ServerSpawnPaintingPacket) { Entity entity = new Entity(entityId, Collections.singletonList(packet)); entity.spawnTime = time; - thingSpawns.put(time, entity); Entity prev = activeEntities.put(entityId, entity); if (prev != null) { - prev.despawnTime = time; + index = prev.writeToCache(indexOut, out, time, index); } } else if (packet instanceof ServerSpawnPlayerPacket) { ServerPlayerListEntryPacket listEntryPacket = new ServerPlayerListEntryPacket( @@ -373,34 +391,30 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe ); Entity entity = new Entity(entityId, Arrays.asList(listEntryPacket, packet)); entity.spawnTime = time; - thingSpawns.put(time, entity); Entity prev = activeEntities.put(entityId, entity); if (prev != null) { - prev.despawnTime = time; + index = prev.writeToCache(indexOut, out, time, index); } } else if (packet instanceof ServerEntityDestroyPacket) { for (int id : ((ServerEntityDestroyPacket) packet).getEntityIds()) { Entity entity = activeEntities.remove(id); if (entity != null) { - entity.despawnTime = time; + index = entity.writeToCache(indexOut, out, time, index); } } } else if (packet instanceof ServerChunkDataPacket) { Column column = ((ServerChunkDataPacket) packet).getColumn(); Chunk chunk = new Chunk(column); chunk.spawnTime = time; - thingSpawns.put(time, chunk); Chunk prev = activeChunks.put(coordToLong(column.getX(), column.getZ()), chunk); if (prev != null) { - prev.currentBlockState = null; // free memory because we no longer need it - prev.despawnTime = time; + index = prev.writeToCache(indexOut, out, time, index); } } else if (packet instanceof ServerUnloadChunkPacket) { ServerUnloadChunkPacket p = (ServerUnloadChunkPacket) packet; Chunk prev = activeChunks.remove(coordToLong(p.getX(), p.getZ())); if (prev != null) { - prev.currentBlockState = null; // free memory because we no longer need it - prev.despawnTime = time; + index = prev.writeToCache(indexOut, out, time, index); } } else if (packet instanceof ServerBlockChangePacket || packet instanceof ServerMultiBlockChangePacket) { for (BlockChangeRecord record : @@ -426,12 +440,16 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } } } else if (packet instanceof ServerRespawnPacket) { - activeEntities.values().forEach(entity -> entity.despawnTime = time); + for (Entity entity : activeEntities.values()) { + index = entity.writeToCache(indexOut, out, time, index); + } activeEntities.clear(); - activeChunks.values().forEach(chunk -> chunk.despawnTime = time); + for (Chunk chunk : activeChunks.values()) { + index = chunk.writeToCache(indexOut, out, time, index); + } activeChunks.clear(); if (activeWeather != null) { - activeWeather.despawnTime = time; + index = activeWeather.writeToCache(indexOut, out, time, index); } activeWeather = null; } else if (packet instanceof ServerUpdateTimePacket) { @@ -441,15 +459,14 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe switch (p.getNotification()) { case START_RAIN: if (activeWeather != null) { - activeWeather.despawnTime = time; + index = activeWeather.writeToCache(indexOut, out, time, index); } activeWeather = new Weather(); activeWeather.spawnTime = time; - thingSpawns.put(time, activeWeather); break; case STOP_RAIN: if (activeWeather != null) { - activeWeather.despawnTime = time; + index = activeWeather.writeToCache(indexOut, out, time, index); activeWeather = null; } break; @@ -475,43 +492,24 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } } } + + for (Entity entity : activeEntities.values()) { + index = entity.writeToCache(indexOut, out, time, index); + } + for (Chunk chunk : activeChunks.values()) { + index = chunk.writeToCache(indexOut, out, time, index); + } + if (activeWeather != null) { + index = activeWeather.writeToCache(indexOut, out, time, index); + } + + indexOut.writeByte(0); + writeToCache(indexOut, worldTimes); + writeToCache(indexOut, thunderStrengths); + + indexOut.writeVarInt(index); } LOGGER.info("Analysed replay in " + (System.currentTimeMillis() - sysTimeStart) + "ms"); - - sysTimeStart = System.currentTimeMillis(); - try (OutputStream cacheOut = replayFile.write(CACHE_ENTRY)) { - NetOutput out = new StreamNetOutput(cacheOut); - out.writeVarInt(CACHE_VERSION); - - out.writeVarInt(thingSpawnsT.size()); - int lastTime = 0; - for (Map.Entry> entry : thingSpawnsT.entrySet()) { - int time = entry.getKey(); - out.writeVarInt(time - lastTime); - lastTime = time; - - progress.accept(time / duration * (1 - progressSplit) + progressSplit); - - Collection trackedThings = entry.getValue(); - out.writeVarInt(trackedThings.size()); - for (TrackedThing trackedThing : trackedThings) { - out.writeVarInt(trackedThing.despawnTime); - if (trackedThing instanceof Entity) { - out.writeVarInt(0); - } else if (trackedThing instanceof Chunk) { - out.writeVarInt(1); - } else if (trackedThing instanceof Weather) { - out.writeVarInt(2); - } else { - throw new UnsupportedOperationException("Unknown type of tracked thing: " + trackedThing); - } - trackedThing.writeToCache(out); - } - } - writeToCache(out, worldTimes); - writeToCache(out, thunderStrengths); - } - LOGGER.info("Wrote quick replay to cache in " + (System.currentTimeMillis() - sysTimeStart) + "ms"); } @Override @@ -519,8 +517,8 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe ensureInitialized(() -> { if (replayTime > currentTimeStamp) { activeThings.removeIf(thing -> { - if (thing.despawnTime < replayTime) { - thing.despawnPackets.forEach(ctx::fireChannelRead); + if (thing.despawnTime <= replayTime) { + thing.despawn(this, ctx::fireChannelRead); return true; } else { return false; @@ -529,7 +527,7 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe thingSpawnsT.subMap(currentTimeStamp, false, replayTime, true).values() .forEach(things -> things.forEach(thing -> { if (thing.despawnTime > replayTime) { - thing.spawnPackets.forEach(ctx::fireChannelRead); + thing.spawn(this, ctx::fireChannelRead); activeThings.add(thing); } })); @@ -539,7 +537,7 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } else { activeThings.removeIf(thing -> { if (thing.spawnTime > replayTime) { - thing.despawnPackets.forEach(ctx::fireChannelRead); + thing.despawn(this, ctx::fireChannelRead); return true; } else { return false; @@ -548,7 +546,7 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe thingDespawnsT.subMap(replayTime, false, currentTimeStamp, true).values() .forEach(things -> things.forEach(thing -> { if (thing.spawnTime <= replayTime) { - thing.spawnPackets.forEach(ctx::fireChannelRead); + thing.spawn(this, ctx::fireChannelRead); activeThings.add(thing); } })); @@ -560,6 +558,25 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe }); } + private static final com.github.steveice10.netty.buffer.ByteBuf mcplByteBuf = com.github.steveice10.netty.buffer.Unpooled.buffer(); + private static final ByteBufNetInput byteBufNetInput = new ByteBufNetInput(mcplByteBuf); + private static final ByteBufNetOutput byteBufNetOutput = new ByteBufNetOutput(mcplByteBuf); + + private static BlockStorage copy(BlockStorage of) { + int readerIndex = byteBuf.readerIndex(); // Mark the current reader and writer index (should be at start) + int writerIndex = byteBuf.writerIndex(); + try { + of.write(byteBufNetOutput); + return new BlockStorage(byteBufNetInput); + } catch (Exception e) { + Utils.throwIfUnchecked(e); + throw new RuntimeException(e); + } finally { + byteBuf.readerIndex(readerIndex); // Reset reader & writer index for next use + byteBuf.writerIndex(writerIndex); + } + } + private static final ByteBuf byteBuf = Unpooled.buffer(); private static final ByteBufOutputStream byteBufOut = new ByteBufOutputStream(byteBuf); private static final PacketBuffer packetBuf = new PacketBuffer(byteBuf); @@ -597,7 +614,7 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } } - private static Packet readPacketFromCache(int version, NetInput in) throws IOException { + private static Packet readPacketFromCache(NetInput in) throws IOException { int readerIndex = byteBuf.readerIndex(); // Mark the current reader and writer index (should be at start) int writerIndex = byteBuf.writerIndex(); try { @@ -622,24 +639,25 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } } - private static List> readPacketsFromCache(int version, NetInput in) throws IOException { + @SneakyThrows(IOException.class) + private static List> readPacketsFromCache(NetInput in) { int size = in.readVarInt(); List> packets = new ArrayList<>(size); for (int i = 0; i < size; i++) { - packets.add(readPacketFromCache(version, in)); + packets.add(readPacketFromCache(in)); } return packets; } - private static void readFromCache(int version, NetInput in, SortedMap> packets) throws IOException { + private static void readFromCache(NetInput in, SortedMap> packets) throws IOException { int time = 0; for (int i = in.readVarInt(); i > 0; i--) { time += in.readVarInt(); - packets.put(time, readPacketFromCache(version, in)); + packets.put(time, readPacketFromCache(in)); } } - private static void writeToCache(NetOutput out, com.github.steveice10.packetlib.packet.Packet packet) throws IOException { + private static int writeToCache(NetOutput out, com.github.steveice10.packetlib.packet.Packet packet) throws IOException { int readerIndex = byteBuf.readerIndex(); // Mark the current reader and writer index (should be at start) int writerIndex = byteBuf.writerIndex(); try { @@ -649,10 +667,11 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe byteBuf.skipBytes(8); // Skip packet length & timestamp int size = byteBuf.readableBytes(); - out.writeVarInt(size); + int len = writeVarInt(out, size); for (int i = 0; i < size; i++) { out.writeByte(byteBuf.readByte()); } + return len + size; } catch (Exception e) { Utils.throwIfInstanceOf(e, IOException.class); Utils.throwIfUnchecked(e); @@ -663,23 +682,37 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } } - private static void writeToCache(NetOutput out, Collection packets) throws IOException { - out.writeVarInt(packets.size()); + private static int writeToCache(NetOutput out, Collection packets) throws IOException { + int len = writeVarInt(out, packets.size()); for (com.github.steveice10.packetlib.packet.Packet packet : packets) { - writeToCache(out, packet); + len += writeToCache(out, packet); } + return len; } - private static void writeToCache(NetOutput out, SortedMap packets) throws IOException { - out.writeVarInt(packets.size()); + private static int writeToCache(NetOutput out, SortedMap packets) throws IOException { + int len = 0; + len += writeVarInt(out, packets.size()); int lastTime = 0; for (Map.Entry entry : packets.entrySet()) { int time = entry.getKey(); - out.writeVarInt(time - lastTime); + len += writeVarInt(out, time - lastTime); lastTime = time; - writeToCache(out, entry.getValue()); + len += writeToCache(out, entry.getValue()); } + return len; + } + + private static int writeVarInt(NetOutput out, int i) throws IOException { + int len = 1; + while ((i & -128) != 0) { + out.writeByte(i & 127 | 128); + i >>>= 7; + len++; + } + out.writeByte(i); + return len; } private static long coordToLong(int x, int z) { @@ -704,7 +737,6 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe List spawnPackets; List despawnPackets; int spawnTime; - int despawnTime = Integer.MAX_VALUE; private TrackedThing(List spawnPackets, List despawnPackets) { @@ -712,26 +744,41 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe this.despawnPackets = despawnPackets; } - public void writeToCache(NetOutput out) throws IOException { - QuickReplaySender.writeToCache(out, spawnPackets); - QuickReplaySender.writeToCache(out, despawnPackets); + public int writeToCache(NetOutput indexOut, NetOutput cacheOut, int despawnTime, int index) throws IOException { + indexOut.writeVarInt(spawnTime); + indexOut.writeVarInt(despawnTime); + indexOut.writeVarInt(index); + index += QuickReplaySender.writeToCache(cacheOut, spawnPackets); + indexOut.writeVarInt(index); + index += QuickReplaySender.writeToCache(cacheOut, despawnPackets); + return index; } } - // For quicker jumping we store MC packets. - // However, during replay analysis it's easier to use MCProtocolLib ones in part because MC packets have been - // renamed from 1.8 to 1.9 but also because we cannot easily serialize them for caching. - // Therefore during analysis we use TrackedThing which we then serialize (we'd have to do that anyway for caching) - // and afterwards unserialize in BakedTrackedThing as MC packets for replaying. + // For memory efficiency we store raw packets (we might even want to consider compression or memory mapped files). + // During analysis we use TrackedThing which we then serialize (we'd have to do that anyway for caching) + // and afterwards deserialize in BakedTrackedThing as MC packets on demand for replaying. private static abstract class BakedTrackedThing { - List> spawnPackets; - List> despawnPackets; + int indexSpawnPackets; + int indexDespawnPackets; int spawnTime; - int despawnTime = Integer.MAX_VALUE; + int despawnTime; - private BakedTrackedThing(int version, NetInput in) throws IOException { - spawnPackets = readPacketsFromCache(version, in); - despawnPackets = readPacketsFromCache(version, in); + private BakedTrackedThing(NetInput in) throws IOException { + spawnTime = in.readVarInt(); + despawnTime = in.readVarInt(); + indexSpawnPackets = in.readVarInt(); + indexDespawnPackets = in.readVarInt(); + } + + public void spawn(QuickReplaySender sender, Consumer> send) { + sender.buf.readerIndex(indexSpawnPackets); + readPacketsFromCache(sender.bufInput).forEach(send); + } + + public void despawn(QuickReplaySender sender, Consumer> send) { + sender.buf.readerIndex(indexDespawnPackets); + readPacketsFromCache(sender.bufInput).forEach(send); } public abstract void play(int currentTimeStamp, int replayTime, Consumer> send); @@ -748,34 +795,53 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } @Override - public void writeToCache(NetOutput out) throws IOException { - super.writeToCache(out); + public int writeToCache(NetOutput indexOut, NetOutput cacheOut, int despawnTime, int index) throws IOException { + indexOut.writeByte(1); + index = super.writeToCache(indexOut, cacheOut, despawnTime, index); - out.writeVarInt(id); - out.writeVarInt(locations.size()); + indexOut.writeVarInt(id); + indexOut.writeVarInt(index); + index += writeVarInt(cacheOut, locations.size()); int lastTime = 0; for (Map.Entry entry : locations.entrySet()) { int time = entry.getKey(); Location loc = entry.getValue(); - out.writeVarInt(time - lastTime); + index += writeVarInt(cacheOut, time - lastTime); lastTime = time; - out.writeDouble(loc.getX()); - out.writeDouble(loc.getY()); - out.writeDouble(loc.getZ()); - out.writeFloat(loc.getYaw()); - out.writeFloat(loc.getPitch()); + cacheOut.writeDouble(loc.getX()); + cacheOut.writeDouble(loc.getY()); + cacheOut.writeDouble(loc.getZ()); + cacheOut.writeFloat(loc.getYaw()); + cacheOut.writeFloat(loc.getPitch()); + index += 32; } + + return index; } } private static class BakedEntity extends BakedTrackedThing { private int id; - private NavigableMap locations = new TreeMap<>(); + private int index; + private NavigableMap locations; - private BakedEntity(int version, NetInput in) throws IOException { - super(version, in); + private BakedEntity(NetInput in) throws IOException { + super(in); id = in.readVarInt(); + index = in.readVarInt(); + } + + @Override + @SneakyThrows(IOException.class) + public void spawn(QuickReplaySender sender, Consumer> send) { + super.spawn(sender, send); + + sender.buf.readerIndex(index); + NetInput in = sender.bufInput; + + locations = new TreeMap<>(); + int time = 0; for (int i = in.readVarInt(); i > 0; i--) { time += in.readVarInt(); @@ -783,6 +849,12 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } } + @Override + public void despawn(QuickReplaySender sender, Consumer> send) { + super.despawn(sender, send); + locations = null; + } + @Override public void play(int currentTimeStamp, int replayTime, Consumer> send) { playMap(locations, currentTimeStamp, replayTime, l -> @@ -806,38 +878,57 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe Collections.singletonList(new ServerUnloadChunkPacket(column.getX(), column.getZ()))); com.github.steveice10.mc.protocol.data.game.chunk.Chunk[] chunks = column.getChunks(); for (int i = 0; i < currentBlockState.length; i++) { - currentBlockState[i] = chunks[i] == null ? new BlockStorage() : chunks[i].getBlocks(); + currentBlockState[i] = chunks[i] == null ? new BlockStorage() : copy(chunks[i].getBlocks()); } } @Override - public void writeToCache(NetOutput out) throws IOException { - super.writeToCache(out); + public int writeToCache(NetOutput indexOut, NetOutput cacheOut, int despawnTime, int index) throws IOException { + indexOut.writeByte(2); + index = super.writeToCache(indexOut, cacheOut, despawnTime, index); - out.writeVarInt(blocksT.size()); + indexOut.writeVarInt(index); + index += writeVarInt(cacheOut, blocksT.size()); int lastTime = 0; for (Map.Entry> entry : blocksT.entrySet()) { int time = entry.getKey(); - out.writeVarInt(time - lastTime); + index += writeVarInt(cacheOut, time - lastTime); lastTime = time; Collection blockChanges = entry.getValue(); - out.writeVarInt(blockChanges.size()); + index += writeVarInt(cacheOut, blockChanges.size()); for (BlockChange blockChange : blockChanges) { - NetUtil.writePosition(out, blockChange.pos); - NetUtil.writeBlockState(out, blockChange.from); - NetUtil.writeBlockState(out, blockChange.to); + NetUtil.writePosition(cacheOut, blockChange.pos); + index += 8; + index += writeVarInt(cacheOut, blockChange.from.getId() << 4 | blockChange.from.getData() & 15); + index += writeVarInt(cacheOut, blockChange.to.getId() << 4 | blockChange.to.getData() & 15); } } + + return index; } } private static class BakedChunk extends BakedTrackedThing { - private TreeMap> blocksT = new TreeMap<>(); - private ListMultimap blocks = Multimaps.newListMultimap(blocksT, LinkedList::new); // LinkedList to allow .descendingIterator + private int index; + private TreeMap> blocksT; - private BakedChunk(int version, NetInput in) throws IOException { - super(version, in); + private BakedChunk(NetInput in) throws IOException { + super(in); + + index = in.readVarInt(); + } + + @Override + @SneakyThrows(IOException.class) + public void spawn(QuickReplaySender sender, Consumer> send) { + super.spawn(sender, send); + + sender.buf.readerIndex(index); + NetInput in = sender.bufInput; + + blocksT = new TreeMap<>(); + ListMultimap blocks = Multimaps.newListMultimap(blocksT, LinkedList::new); // LinkedList to allow .descendingIterator int time = 0; for (int i = in.readVarInt(); i > 0; i--) { @@ -853,6 +944,13 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } } + @Override + public void despawn(QuickReplaySender sender, Consumer> send) { + super.despawn(sender, send); + + blocksT = null; + } + @Override public void play(int currentTimeStamp, int replayTime, Consumer> send) { blocksT.subMap(currentTimeStamp, false, replayTime, true).values() @@ -895,20 +993,44 @@ public class QuickReplaySender extends ChannelHandlerAdapter implements ReplaySe } @Override - public void writeToCache(NetOutput out) throws IOException { - super.writeToCache(out); + public int writeToCache(NetOutput indexOut, NetOutput cacheOut, int despawnTime, int index) throws IOException { + indexOut.writeByte(3); + index = super.writeToCache(indexOut, cacheOut, despawnTime, index); - QuickReplaySender.writeToCache(out, rainStrengths); + indexOut.writeVarInt(index); + index += QuickReplaySender.writeToCache(cacheOut, rainStrengths); + + return index; } } private static class BakedWeather extends BakedTrackedThing { - private TreeMap> rainStrengths = new TreeMap<>(); + private int index; + private TreeMap> rainStrengths; - private BakedWeather(int version, NetInput in) throws IOException { - super(version, in); + private BakedWeather(NetInput in) throws IOException { + super(in); - readFromCache(version, in, rainStrengths); + index = in.readVarInt(); + } + + @Override + @SneakyThrows(IOException.class) + public void spawn(QuickReplaySender sender, Consumer> send) { + super.spawn(sender, send); + + sender.buf.readerIndex(index); + + rainStrengths = new TreeMap<>(); + + readFromCache(sender.bufInput, rainStrengths); + } + + @Override + public void despawn(QuickReplaySender sender, Consumer> send) { + super.despawn(sender, send); + + rainStrengths = null; } @Override