- private TransmittedConnectionEntry transmit(final ConnectionEntry entry) {
- final long txSequence = nextTxSequence++;
-
- final RequestEnvelope toSend = new RequestEnvelope(entry.getRequest().toVersion(remoteVersion()), sessionId(),
- txSequence);
-
- final ActorRef actor = remoteActor();
- LOG.trace("Transmitting request {} as {} to {}", entry.getRequest(), toSend, actor);
- actor.tell(toSend, ActorRef.noSender());
-
- return new TransmittedConnectionEntry(entry, sessionId(), txSequence, readTime());
- }
-
- @Override
- void enqueueEntry(final ConnectionEntry entry) {
- if (inflightSize() < remoteMaxMessages()) {
- appendToInflight(transmit(entry));
- LOG.debug("Enqueued request {} to queue {}", entry.getRequest(), this);
- } else {
- LOG.debug("Queue is at capacity, delayed sending of request {}", entry.getRequest());
- super.enqueueEntry(entry);
- }
- }
-
- @Override
- void sendMessages(final int count) {
- int toSend = count;
-
- while (toSend > 0) {
- final ConnectionEntry e = dequeEntry();
- if (e == null) {
- break;
- }
-
- LOG.debug("Transmitting entry {}", e);
- appendToInflight(transmit(e));
- toSend--;
- }
- }
-