import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions;
import com.google.common.base.Verify;
-import com.google.common.collect.Iterables;
import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.ArrayDeque;
+import java.util.Collection;
import java.util.Deque;
import java.util.Iterator;
import java.util.Optional;
tracker = new AveragingProgressTracker(targetDepth);
}
- final Iterable<ConnectionEntry> asIterable() {
- return Iterables.concat(inflight, pending);
+ /**
+ * Drain the contents of the connection into a list. This will leave the queue empty and allow further entries
+ * to be added to it during replay. When we set the successor all entries enqueued between when this methods
+ * returns and the successor is set will be replayed to the successor.
+ *
+ * @return Collection of entries present in the queue.
+ */
+ final Collection<ConnectionEntry> drain() {
+ final Collection<ConnectionEntry> ret = new ArrayDeque<>(inflight.size() + pending.size());
+ ret.addAll(inflight);
+ ret.addAll(pending);
+ inflight.clear();
+ pending.clear();
+ return ret;
}
final long ticksStalling(final long now) {
tracker.closeTask(now, entry.getEnqueuedTicks(), entry.getTxTicks(), envelope.getExecutionTimeNanos());
// We have freed up a slot, try to transmit something
+ tryTransmit(now);
+
+ return Optional.of(entry);
+ }
+
+ final void tryTransmit(final long now) {
final int toSend = canTransmitCount(inflight.size());
if (toSend > 0 && !pending.isEmpty()) {
transmitEntries(toSend, now);
}
-
- return Optional.of(entry);
}
private void transmitEntries(final int maxTransmit, final long now) {
*/
final long enqueue(final ConnectionEntry entry, final long now) {
if (successor != null) {
+ // This call will pay the enqueuing price, hence the caller does not have to
successor.forwardEntry(entry, now);
return 0;
}
successor = Preconditions.checkNotNull(forwarder);
LOG.debug("Connection {} superseded by {}, splicing queue", this, successor);
+ /*
+ * We need to account for entries which have been added between the time drain() was called and this method
+ * is invoked. Since the old connection is visible during replay and some entries may have completed on the
+ * replay thread, there was an avenue for this to happen.
+ */
+ int count = 0;
ConnectionEntry entry = inflight.poll();
while (entry != null) {
- successor.forwardEntry(entry, now);
+ successor.replayEntry(entry, now);
entry = inflight.poll();
+ count++;
}
entry = pending.poll();
while (entry != null) {
- successor.forwardEntry(entry, now);
+ successor.replayEntry(entry, now);
entry = pending.poll();
+ count++;
+ }
+
+ LOG.debug("Connection {} queue spliced {} messages", this, count);
+ }
+
+ final void remove(final long now) {
+ final TransmittedConnectionEntry txe = inflight.poll();
+ if (txe == null) {
+ final ConnectionEntry entry = pending.pop();
+ tracker.closeTask(now, entry.getEnqueuedTicks(), 0, 0);
+ } else {
+ tracker.closeTask(now, txe.getEnqueuedTicks(), txe.getTxTicks(), 0);
}
}
}
queue.clear();
}
-
}