=act #13861 report cause for Akka IO TCP CommandFailed events
To stay compatible this needs to be hacked into CommandFailed using a mutable var.
This commit is contained in:
parent
2858118946
commit
38c9fff219
7 changed files with 73 additions and 28 deletions
|
|
@ -4,21 +4,22 @@
|
|||
|
||||
package akka.io
|
||||
|
||||
import java.net.{ SocketException, InetSocketAddress }
|
||||
import java.net.{ InetSocketAddress, SocketException }
|
||||
import java.nio.channels.SelectionKey._
|
||||
import java.io.{ FileInputStream, IOException }
|
||||
import java.nio.channels.{ FileChannel, SocketChannel }
|
||||
import java.nio.ByteBuffer
|
||||
|
||||
import scala.annotation.tailrec
|
||||
import scala.collection.immutable
|
||||
import scala.util.control.NonFatal
|
||||
import scala.util.control.{ NoStackTrace, NonFatal }
|
||||
import scala.concurrent.duration._
|
||||
import akka.actor._
|
||||
import akka.util.ByteString
|
||||
import akka.io.Inet.SocketOption
|
||||
import akka.io.Tcp._
|
||||
import akka.io.SelectionHandler._
|
||||
import akka.dispatch.{ UnboundedMessageQueueSemantics, RequiresMessageQueue }
|
||||
import akka.dispatch.{ RequiresMessageQueue, UnboundedMessageQueueSemantics }
|
||||
import java.nio.file.Paths
|
||||
|
||||
/**
|
||||
|
|
@ -150,11 +151,11 @@ private[io] abstract class TcpConnection(val tcp: TcpExt, val channel: SocketCha
|
|||
case write: WriteCommand ⇒
|
||||
if (writingSuspended) {
|
||||
if (TraceLogging) log.debug("Dropping write because writing is suspended")
|
||||
sender() ! write.failureMessage
|
||||
sender() ! write.failureMessage.withCause(DroppingWriteBecauseWritingIsSuspendedException)
|
||||
|
||||
} else if (writePending) {
|
||||
if (TraceLogging) log.debug("Dropping write because queue is full")
|
||||
sender() ! write.failureMessage
|
||||
sender() ! write.failureMessage.withCause(DroppingWriteBecauseQueueIsFullException)
|
||||
if (info.useResumeWriting) writingSuspended = true
|
||||
|
||||
} else {
|
||||
|
|
@ -505,4 +506,10 @@ private[io] object TcpConnection {
|
|||
}
|
||||
|
||||
val doNothing: () ⇒ Unit = () ⇒ ()
|
||||
|
||||
val DroppingWriteBecauseWritingIsSuspendedException =
|
||||
new IOException("Dropping write because writing is suspended") with NoStackTrace
|
||||
|
||||
val DroppingWriteBecauseQueueIsFullException =
|
||||
new IOException("Dropping write because queue is full") with NoStackTrace
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue