mirror of
https://github.com/v2fly/v2ray-core.git
synced 2024-12-31 14:36:50 -05:00
check connection state for every write operation
This commit is contained in:
parent
b5caea67ac
commit
a6c0ef11ba
@ -340,22 +340,31 @@ func (c *Connection) waitForDataOutput() error {
|
||||
func (c *Connection) Write(b []byte) (int, error) {
|
||||
totalWritten := 0
|
||||
|
||||
for {
|
||||
dataWritten := false
|
||||
for {
|
||||
if c == nil || c.State() != StateActive {
|
||||
return totalWritten, io.ErrClosedPipe
|
||||
}
|
||||
|
||||
for c.sendingWorker.Push(func(bb []byte) (int, error) {
|
||||
if !c.sendingWorker.Push(func(bb []byte) (int, error) {
|
||||
n := copy(bb[:c.mss], b[totalWritten:])
|
||||
totalWritten += n
|
||||
return n, nil
|
||||
}) {
|
||||
c.dataUpdater.WakeUp()
|
||||
break
|
||||
}
|
||||
|
||||
dataWritten = true
|
||||
|
||||
if totalWritten == len(b) {
|
||||
return totalWritten, nil
|
||||
}
|
||||
}
|
||||
|
||||
if dataWritten {
|
||||
c.dataUpdater.WakeUp()
|
||||
}
|
||||
|
||||
if err := c.waitForDataOutput(); err != nil {
|
||||
return totalWritten, err
|
||||
}
|
||||
@ -366,20 +375,28 @@ func (c *Connection) Write(b []byte) (int, error) {
|
||||
func (c *Connection) WriteMultiBuffer(mb buf.MultiBuffer) error {
|
||||
defer mb.Release()
|
||||
|
||||
for {
|
||||
dataWritten := false
|
||||
for {
|
||||
if c == nil || c.State() != StateActive {
|
||||
return io.ErrClosedPipe
|
||||
}
|
||||
|
||||
for c.sendingWorker.Push(func(bb []byte) (int, error) {
|
||||
if !c.sendingWorker.Push(func(bb []byte) (int, error) {
|
||||
return mb.Read(bb[:c.mss])
|
||||
}) {
|
||||
c.dataUpdater.WakeUp()
|
||||
break
|
||||
}
|
||||
dataWritten = true
|
||||
if mb.IsEmpty() {
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
if dataWritten {
|
||||
c.dataUpdater.WakeUp()
|
||||
}
|
||||
|
||||
if err := c.waitForDataOutput(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
Loading…
Reference in New Issue
Block a user