Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import com.comphenix.protocol.events.PacketListener;
import com.comphenix.protocol.injector.collection.InboundPacketListenerSet;
import com.comphenix.protocol.injector.collection.OutboundPacketListenerSet;
import com.comphenix.protocol.injector.netty.Injector;
import com.comphenix.protocol.scheduler.ProtocolScheduler;
import com.google.common.base.Objects;
import com.google.common.collect.ImmutableList;
Expand Down Expand Up @@ -274,7 +275,7 @@ void unregisterAsyncHandlerInternal(AsyncListenerHandler handler) {
* @return TRUE if we are, FALSE otherwise.
*/
private boolean onMainThread() {
return Thread.currentThread().getId() == mainThread.getId();
return Thread.currentThread() == mainThread;
}

@Override
Expand Down Expand Up @@ -347,22 +348,22 @@ public boolean hasAsynchronousListeners(PacketEvent packet) {
* Construct a asynchronous marker with all the default values.
* @return Asynchronous marker.
*/
public AsyncMarker createAsyncMarker() {
return createAsyncMarker(AsyncMarker.DEFAULT_TIMEOUT_DELTA);
public AsyncMarker createAsyncMarker(Injector injector) {
return createAsyncMarker(injector, AsyncMarker.DEFAULT_TIMEOUT_DELTA);
}

/**
* Construct an async marker with the given sending priority delta and timeout delta.
* @param timeoutDelta - how long (in ms) until the packet expire.
* @return An async marker.
*/
public AsyncMarker createAsyncMarker(long timeoutDelta) {
return createAsyncMarker(timeoutDelta, currentSendingIndex.incrementAndGet());
public AsyncMarker createAsyncMarker(Injector injector, long timeoutDelta) {
return createAsyncMarker(injector, timeoutDelta, currentSendingIndex.incrementAndGet());
}

// Helper method
private AsyncMarker createAsyncMarker(long timeoutDelta, long sendingIndex) {
return new AsyncMarker(manager, sendingIndex, System.currentTimeMillis(), timeoutDelta);
private AsyncMarker createAsyncMarker(Injector injector, long timeoutDelta, long sendingIndex) {
return new AsyncMarker(injector, sendingIndex, System.currentTimeMillis(), timeoutDelta);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -538,7 +538,7 @@ private boolean waitForStops() throws InterruptedException {
*/
private void listenerLoop(int workerID) {
// Danger, danger!
if (Thread.currentThread().getId() == mainThread.getId())
if (Thread.currentThread() == mainThread)
throw new IllegalStateException("Do not call this method from the main thread.");
if (cancelled)
throw new IllegalStateException("Listener has been cancelled. Create a new listener instead.");
Expand Down
36 changes: 22 additions & 14 deletions src/main/java/com/comphenix/protocol/async/AsyncMarker.java
Original file line number Diff line number Diff line change
Expand Up @@ -23,14 +23,18 @@
import java.lang.reflect.Method;
import java.util.Iterator;
import java.util.List;
import java.util.Objects;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.logging.Level;

import com.comphenix.protocol.PacketStream;
import com.comphenix.protocol.PacketType;
import com.comphenix.protocol.ProtocolLibrary;
import com.comphenix.protocol.ProtocolLogger;
import com.comphenix.protocol.events.NetworkMarker;
import com.comphenix.protocol.events.PacketEvent;
import com.comphenix.protocol.injector.netty.Injector;
import com.comphenix.protocol.reflect.FieldAccessException;
import com.comphenix.protocol.reflect.FuzzyReflection;
import com.comphenix.protocol.utility.MinecraftReflection;
Expand Down Expand Up @@ -64,7 +68,7 @@ public class AsyncMarker implements Serializable, Comparable<AsyncMarker> {
/**
* The packet stream responsible for transmitting the packet when it's done processing.
*/
private transient PacketStream packetStream;
private final Injector injector;

/**
* Current list of async packet listeners.
Expand All @@ -86,7 +90,7 @@ public class AsyncMarker implements Serializable, Comparable<AsyncMarker> {
private volatile boolean processed;

// Whether or not the packet has been sent
private volatile boolean transmitted;
private final AtomicBoolean transmitted = new AtomicBoolean(false);

// Whether or not the asynchronous processing itself should be cancelled
private volatile boolean asyncCancelled;
Expand All @@ -109,11 +113,8 @@ public class AsyncMarker implements Serializable, Comparable<AsyncMarker> {
* Create a container for asyncronous packets.
* @param initialTime - the current time in milliseconds since 01.01.1970 00:00.
*/
AsyncMarker(PacketStream packetStream, long sendingIndex, long initialTime, long timeoutDelta) {
if (packetStream == null)
throw new IllegalArgumentException("packetStream cannot be NULL");

this.packetStream = packetStream;
AsyncMarker(Injector injector, long sendingIndex, long initialTime, long timeoutDelta) {
this.injector = Objects.requireNonNull(injector, "injector is nul");

// Timeout
this.initialTime = initialTime;
Expand Down Expand Up @@ -180,16 +181,18 @@ public void setNewSendingIndex(long newSendingIndex) {
* Retrieve the packet stream responsible for transmitting this packet.
* @return The packet stream.
*/
@Deprecated
public PacketStream getPacketStream() {
return packetStream;
return ProtocolLibrary.getProtocolManager();
}

/**
* Sets the output packet stream responsible for transmitting this packet.
* @param packetStream - new output packet stream.
*/
@Deprecated
public void setPacketStream(PacketStream packetStream) {
this.packetStream = packetStream;
// NOOP
}

/**
Expand Down Expand Up @@ -285,7 +288,7 @@ public void setProcessingLock(Object processingLock) {
* @return TRUE if it has been sent before, FALSE otherwise.
*/
public boolean isTransmitted() {
return transmitted;
return transmitted.get();
}

/**
Expand Down Expand Up @@ -383,13 +386,18 @@ void setListenerTraversal(Iterator<AsyncListenerHandler> listenerTraversal) {
* @throws IOException If the packet couldn't be sent.
*/
void sendPacket(PacketEvent event) throws IOException {
if (!this.transmitted.compareAndSet(false, true)) {
return;
}

Object handle = event.getPacket().getHandle();

if (event.isServerPacket()) {
packetStream.sendServerPacket(event.getPlayer(), event.getPacket(), NetworkMarker.getNetworkMarker(event), false);
NetworkMarker marker = NetworkMarker.getNetworkMarker(event);
this.injector.sendClientboundPacket(handle, marker, false);
} else {
packetStream.receiveClientPacket(event.getPlayer(), event.getPacket(), NetworkMarker.getNetworkMarker(event),
false);
this.injector.readServerboundPacket(handle);
}
transmitted = true;
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
import com.comphenix.protocol.injector.collection.InboundPacketListenerSet;
import com.comphenix.protocol.injector.collection.OutboundPacketListenerSet;
import com.comphenix.protocol.injector.collection.PacketListenerSet;
import com.comphenix.protocol.injector.netty.Injector;
import com.comphenix.protocol.injector.netty.WirePacket;
import com.comphenix.protocol.injector.netty.manager.NetworkManagerInjector;
import com.comphenix.protocol.injector.packet.PacketRegistry;
Expand Down Expand Up @@ -151,29 +152,20 @@ public void sendServerPacket(Player receiver, PacketContainer packet, boolean fi

@Override
public void sendServerPacket(Player receiver, PacketContainer packet, NetworkMarker marker, boolean filters) {
if (!this.closed) {
// if we skip the packet events later when actually writing into the pipeline we at least notify all
// monitor listeners before doing so - they will not be able to change the event tho
if (!filters) {
// ensure we are on the main thread if any listener requires that
if (this.hasMainThreadListener(packet.getType()) && !this.server.isPrimaryThread()) {
NetworkMarker copy = marker; // okay fine
ProtocolLibrary.getScheduler().scheduleSyncDelayedTask(
() -> this.sendServerPacket(receiver, packet, copy, false), 1L);
return;
}

// construct the event and post to all monitor listeners
if (this.closed) {
return;
}

// A monitor listener should never modify a packet/event so we can simply invoke monitor listeners
// independently of our injector pipeline
if (!filters) {
this.runMonitorListeners(packet, () -> {
PacketEvent event = PacketEvent.fromServer(this, packet, marker, receiver, false);
this.outboundListeners.invoke(event, ListenerPriority.MONITOR);

// update the marker of the event without accidentally constructing it
marker = NetworkMarker.getNetworkMarker(event);
}

// process outbound
this.networkManagerInjector.getInjector(receiver).sendClientboundPacket(packet.getHandle(), marker, filters);
});
}

this.networkManagerInjector.getInjector(receiver).sendClientboundPacket(packet.getHandle(), marker, filters);
}

@Override
Expand All @@ -200,34 +192,40 @@ public void receiveClientPacket(Player sender, PacketContainer packet, boolean f

@Override
public void receiveClientPacket(Player sender, PacketContainer packet, NetworkMarker marker, boolean filters) {
if (!this.closed) {
// make sure we are on the main thread if any listener of the packet needs it
if (this.hasMainThreadListener(packet.getType()) && !this.server.isPrimaryThread()) {
ProtocolLibrary.getScheduler().runTask(
() -> this.receiveClientPacket(sender, packet, marker, filters));
if (this.closed) {
return;
}

// make sure we are on the main thread if any listener of the packet needs it
if (filters && this.requiresMainThread(packet)) {
ProtocolLibrary.getScheduler().runTask(
() -> this.receiveClientPacket(sender, packet, marker, filters));
return;
}

Object nmsPacket = packet.getHandle();

if (filters) {
PacketEvent event = PacketEvent.fromClient(this, packet, marker, sender, false);
this.invokeInboundPacketListeners(event);

if (event.isCancelled()) {
return;
}

Object nmsPacket = packet.getHandle();
// check to which listeners we need to post the packet
if (filters) {
// post to all listeners
PacketEvent event = PacketEvent.fromClient(this.networkManagerInjector, packet, null, sender);
this.invokeInboundPacketListeners(event);
if (event.isCancelled()) {
return;
}

// prevent possible de-sync
nmsPacket = event.getPacket().getHandle();
} else {
// Prevent possible de-sync if the packet was replaced by a listener.
nmsPacket = event.getPacket().getHandle();
} else {
// A monitor listener should never modify a packet/event so we can simply invoke monitor listeners
// independently of our injector pipeline
this.runMonitorListeners(packet, () -> {
PacketEvent event = PacketEvent.fromClient(this, packet, marker, sender, false);
this.inboundListeners.invoke(event, ListenerPriority.MONITOR);
}

// post to the player inject, reset our cancel state change
this.networkManagerInjector.getInjector(sender).readServerboundPacket(nmsPacket);
});
}

// post to the player inject, reset our cancel state change
this.networkManagerInjector.getInjector(sender).readServerboundPacket(nmsPacket);
}

@Override
Expand Down Expand Up @@ -516,12 +514,25 @@ public void invokeOutboundPacketListeners(PacketEvent event) {
this.postPacketToListeners(this.outboundListeners, event, true);
}
}

private boolean requiresMainThread(PacketContainer packet) {
return this.hasMainThreadListener(packet.getType()) && !this.server.isPrimaryThread();
}

private void runMonitorListeners(PacketContainer packet, Runnable notifyMonitor) {
if (this.requiresMainThread(packet)) {
ProtocolLibrary.getScheduler().runTask(notifyMonitor);
} else {
notifyMonitor.run();
}
}

private void postPacketToListeners(PacketListenerSet listeners, PacketEvent event, boolean outbound) {
try {
// append async marker if any async listener for the packet was registered
if (this.asyncFilterManager.hasAsynchronousListeners(event)) {
event.setAsyncMarker(this.asyncFilterManager.createAsyncMarker());
Injector injector = this.networkManagerInjector.getInjector(event.getPlayer());
event.setAsyncMarker(this.asyncFilterManager.createAsyncMarker(injector));
}

// post to sync listeners
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -243,7 +243,7 @@ public void close() {
this.uninject();

// remove any outgoing references
this.channel.attr(INJECTOR).remove();
this.channel.attr(INJECTOR).set(null);

// cleanup
this.savedMarkers.clear();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -267,6 +267,7 @@ public ChannelFuture register(ChannelPromise channelPromise) {
return this.delegate.register(channelPromise);
}

@Deprecated
@Override
public ChannelFuture register(Channel channel, ChannelPromise promise) {
return this.delegate.register(channel, promise);
Expand Down
Loading