package quic
import (
"context"
"fmt"
"io"
"net"
"sync"
"time"
"github.com/tmc/go-iroh/internal/qng/internal/ackhandler"
"github.com/tmc/go-iroh/internal/qng/internal/flowcontrol"
"github.com/tmc/go-iroh/internal/qng/internal/monotime"
"github.com/tmc/go-iroh/internal/qng/internal/protocol"
"github.com/tmc/go-iroh/internal/qng/internal/wire"
)
type SendStream struct {
mutex sync .Mutex
numOutstandingFrames int64
retransmissionQueue []*wire .StreamFrame
ctx context .Context
ctxCancel context .CancelCauseFunc
streamID protocol .StreamID
sender streamSender
reliableSize protocol .ByteCount
writeOffset protocol .ByteCount
shutdownErr error
resetErr *StreamError
queuedResetStreamFrame *wire .ResetStreamFrame
supportsResetStreamAt bool
finishedWriting bool
finSent bool
cancellationFlagged bool
completed bool
dataForWriting []byte
writeBuffer []byte
writeBufferHead int
writeBufferLimit int
active bool
writesInEpisode uint16
burstUntil monotime .Time
corkPending bool
activationTimer *time .Timer
activationGen uint64
writeChan chan struct {}
writeActive bool
writeWake chan struct {}
deadline monotime .Time
flowController flowcontrol .StreamFlowController
}
const (
sendStreamWriteBufferSize = 4096
sendStreamWriteBufferMaxSize = 65536
maxBufferedWriteSize = sendStreamWriteBufferMaxSize / 4
)
var (
_ io .ReaderFrom = &SendStream {}
_ streamControlFrameGetter = &SendStream {}
_ outgoingStream = &SendStream {}
_ sendStreamFrameHandler = &SendStream {}
)
func newSendStream(
ctx context .Context ,
streamID protocol .StreamID ,
sender streamSender ,
flowController flowcontrol .StreamFlowController ,
supportsResetStreamAt bool ,
) *SendStream {
s := &SendStream {
streamID : streamID ,
sender : sender ,
flowController : flowController ,
writeChan : make (chan struct {}, 1 ),
writeWake : make (chan struct {}, 1 ),
supportsResetStreamAt : supportsResetStreamAt ,
}
s .ctx , s .ctxCancel = context .WithCancelCause (ctx )
return s
}
func (s *SendStream ) StreamID () StreamID {
return s .streamID
}
func (s *SendStream ) Write (p []byte ) (int , error ) {
s .mutex .Lock ()
if !s .writeActive && s .resetErr == nil && s .shutdownErr == nil &&
!s .finishedWriting && s .deadline .IsZero () && len (p ) > 0 &&
len (p ) <= maxBufferedWriteSize &&
s .active && s .growWriteBufferFor (len (p )) {
s .appendWriteBuffer (p )
if s .writesInEpisode < ^uint16 (0 ) {
s .writesInEpisode ++
}
s .mutex .Unlock ()
return len (p ), nil
}
for s .writeActive {
s .mutex .Unlock ()
<-s .writeWake
s .mutex .Lock ()
}
s .writeActive = true
isNewlyCompleted , n , err := s .writeLocked (p )
if isNewlyCompleted {
s .sender .onStreamCompleted (s .streamID )
}
return n , err
}
func (s *SendStream ) ReadFrom (r io .Reader ) (int64 , error ) {
var total int64
buf := make ([]byte , sendStreamWriteBufferSize )
for {
n , rerr := r .Read (buf )
if n > 0 {
wn , werr := s .Write (buf [:n ])
total += int64 (wn )
if werr != nil {
return total , werr
}
}
if rerr == io .EOF {
return total , nil
}
if rerr != nil {
return total , rerr
}
}
}
func (s *SendStream ) Writev (bufs *net .Buffers ) (int64 , error ) {
total , err := s .writeVectored (*bufs )
consumed := total
for consumed > 0 && len (*bufs ) > 0 {
if n := int64 (len ((*bufs )[0 ])); consumed >= n {
consumed -= n
*bufs = (*bufs )[1 :]
continue
}
(*bufs )[0 ] = (*bufs )[0 ][consumed :]
consumed = 0
}
if total > 0 && len (*bufs ) > 0 && err == nil {
err = io .ErrShortWrite
}
return total , err
}
func (s *SendStream ) writeVectored (bufs [][]byte ) (int64 , error ) {
var total int64
s .mutex .Lock ()
appended := false
for i := 0 ; i < len (bufs ); {
p := bufs [i ]
if len (p ) == 0 {
i ++
continue
}
if !s .writeActive && s .resetErr == nil && s .shutdownErr == nil &&
!s .finishedWriting && s .deadline .IsZero () &&
s .active && len (p ) <= maxBufferedWriteSize && s .growWriteBufferFor (len (p )) {
s .appendWriteBuffer (p )
appended = true
total += int64 (len (p ))
i ++
continue
}
s .mutex .Unlock ()
n , err := s .Write (p )
total += int64 (n )
if err != nil {
return total , err
}
i ++
s .mutex .Lock ()
}
if appended && s .writesInEpisode < ^uint16 (0 ) {
s .writesInEpisode ++
}
s .mutex .Unlock ()
return total , nil
}
func (s *SendStream ) writeLocked (p []byte ) (bool , int , error ) {
if s .resetErr != nil {
s .cancellationFlagged = true
completed := s .isNewlyCompleted ()
err := s .resetErr
s .finishWriteLocked ()
return completed , 0 , err
}
if s .shutdownErr != nil {
err := s .shutdownErr
s .finishWriteLocked ()
return false , 0 , err
}
if s .finishedWriting {
s .finishWriteLocked ()
return false , 0 , fmt .Errorf ("write on closed stream %d" , s .streamID )
}
if !s .deadline .IsZero () && !monotime .Now ().Before (s .deadline ) {
s .finishWriteLocked ()
return false , 0 , errDeadline
}
if len (p ) == 0 {
s .finishWriteLocked ()
return false , 0 , nil
}
s .dataForWriting = p
var (
deadlineTimer *time .Timer
bytesWritten int
notifiedSender bool
)
for {
var copied bool
var deadline monotime .Time
if s .shutdownErr != nil || s .resetErr != nil {
break
}
canBuffer := s .canBufferWrite ()
if !canBuffer && len (p ) <= maxBufferedWriteSize {
canBuffer = s .growWriteBufferFor (len (s .dataForWriting ))
}
if canBuffer && len (s .dataForWriting ) > 0 {
s .appendWriteBuffer (s .dataForWriting )
s .dataForWriting = nil
bytesWritten = len (p )
copied = true
if s .writesInEpisode < ^uint16 (0 ) {
s .writesInEpisode ++
}
} else {
bytesWritten = len (p ) - len (s .dataForWriting )
deadline = s .deadline
if !deadline .IsZero () {
if !monotime .Now ().Before (deadline ) {
s .dataForWriting = nil
s .finishWriteLocked ()
return false , bytesWritten , errDeadline
}
if deadlineTimer == nil {
deadlineTimer = time .NewTimer (monotime .Until (deadline ))
defer deadlineTimer .Stop ()
} else {
deadlineTimer .Reset (monotime .Until (deadline ))
}
}
if s .dataForWriting == nil || s .shutdownErr != nil || s .resetErr != nil {
break
}
}
notifySender := false
if !notifiedSender {
notifiedSender = true
if !s .active {
notifySender = s .activateOrDelayLocked ()
}
}
if copied {
s .writeActive = false
}
s .mutex .Unlock ()
if notifySender {
s .sender .onHasStreamData (s .streamID , s )
}
if copied {
s .wakeWriter ()
return false , bytesWritten , nil
}
if deadline .IsZero () {
<-s .writeChan
} else {
select {
case <- s .writeChan :
case <- deadlineTimer .C :
}
}
s .mutex .Lock ()
}
if bytesWritten == len (p ) {
s .finishWriteLocked ()
return false , bytesWritten , nil
}
if s .shutdownErr != nil {
err := s .shutdownErr
s .finishWriteLocked ()
return false , bytesWritten , err
}
if s .resetErr != nil {
s .cancellationFlagged = true
completed := s .isNewlyCompleted ()
err := s .resetErr
s .finishWriteLocked ()
return completed , bytesWritten , err
}
s .finishWriteLocked ()
return false , bytesWritten , nil
}
func (s *SendStream ) finishWriteLocked () {
s .writeActive = false
s .mutex .Unlock ()
s .wakeWriter ()
}
func (s *SendStream ) wakeWriter () {
select {
case s .writeWake <- struct {}{}:
default :
}
}
func (s *SendStream ) bufferedWriteLen () int {
return len (s .writeBuffer ) - s .writeBufferHead
}
func (s *SendStream ) writeBufferLimitLocked () int {
if s .writeBufferLimit == 0 {
return sendStreamWriteBufferSize
}
return s .writeBufferLimit
}
func (s *SendStream ) growWriteBufferFor (n int ) bool {
limit := s .writeBufferLimitLocked ()
need := s .bufferedWriteLen () + n
for need > limit && limit < sendStreamWriteBufferMaxSize {
limit *= 2
}
s .writeBufferLimit = limit
return need <= limit
}
func (s *SendStream ) canBufferWrite () bool {
return s .bufferedWriteLen ()+len (s .dataForWriting ) <= s .writeBufferLimitLocked ()
}
func (s *SendStream ) appendWriteBuffer (p []byte ) {
if s .writeBuffer == nil {
s .writeBuffer = make ([]byte , 0 , sendStreamWriteBufferSize )
}
if cap (s .writeBuffer )-len (s .writeBuffer ) < len (p ) {
copy (s .writeBuffer , s .writeBuffer [s .writeBufferHead :])
s .writeBuffer = s .writeBuffer [:s .bufferedWriteLen ()]
s .writeBufferHead = 0
}
s .writeBuffer = append (s .writeBuffer , p ...)
}
func (s *SendStream ) activateOrDelayLocked () bool {
if !s .corkPending {
if s .burstUntil .IsZero () || !monotime .Now ().Before (s .burstUntil ) {
s .burstUntil = 0
s .active = true
return true
}
s .corkPending = true
}
if sendStreamTailDelay <= 0 || s .bufferedWriteLen () >= sendStreamActivationThreshold {
s .stopActivationTimerLocked ()
s .burstUntil = 0
s .corkPending = false
s .active = true
return true
}
if s .activationTimer == nil {
s .activationGen ++
gen := s .activationGen
s .activationTimer = time .AfterFunc (sendStreamTailDelay , func () {
s .activateAfterDelay (gen )
})
}
return false
}
func (s *SendStream ) stopActivationTimerLocked () {
if s .activationTimer == nil {
return
}
s .activationTimer .Stop ()
s .activationTimer = nil
s .activationGen ++
}
func (s *SendStream ) activateAfterDelay (gen uint64 ) {
s .mutex .Lock ()
if gen != s .activationGen {
s .mutex .Unlock ()
return
}
s .activationTimer = nil
if s .active || s .shutdownErr != nil || s .resetErr != nil ||
(s .bufferedWriteLen () == 0 && s .dataForWriting == nil ) {
s .mutex .Unlock ()
return
}
s .burstUntil = 0
s .corkPending = false
s .active = true
s .mutex .Unlock ()
recordCorkTimerActivation ()
s .sender .onHasStreamData (s .streamID , s )
}
func (s *SendStream ) popStreamFrame (maxBytes protocol .ByteCount , v protocol .Version ) (_ ackhandler .StreamFrame , _ *wire .StreamDataBlockedFrame , hasMore bool ) {
s .mutex .Lock ()
f , blocked , hasMoreData := s .popNewOrRetransmittedStreamFrame (maxBytes , v )
if f != nil {
s .numOutstandingFrames ++
}
if !hasMoreData {
if s .writesInEpisode >= sendStreamBurstMinWrites {
s .burstUntil = monotime .Now ().Add (sendStreamBurstFreshness )
} else {
s .burstUntil = 0
}
s .writesInEpisode = 0
s .corkPending = false
s .active = false
}
s .mutex .Unlock ()
if f == nil {
return ackhandler .StreamFrame {}, blocked , hasMoreData
}
return ackhandler .StreamFrame {
Frame : f ,
Handler : (*sendStreamAckHandler )(s ),
}, blocked , hasMoreData
}
func (s *SendStream ) notifyHasStreamData () {
s .mutex .Lock ()
s .stopActivationTimerLocked ()
s .burstUntil = 0
s .corkPending = false
if s .active {
s .mutex .Unlock ()
return
}
s .active = true
s .mutex .Unlock ()
s .sender .onHasStreamData (s .streamID , s )
}
func (s *SendStream ) popNewOrRetransmittedStreamFrame (maxBytes protocol .ByteCount , v protocol .Version ) (_ *wire .StreamFrame , _ *wire .StreamDataBlockedFrame , hasMoreData bool ) {
if s .shutdownErr != nil {
return nil , nil , false
}
if s .resetErr != nil {
reliableOffset := s .reliableOffset ()
if reliableOffset == 0 || (s .writeOffset >= reliableOffset && len (s .retransmissionQueue ) == 0 ) {
return nil , nil , false
}
}
if len (s .retransmissionQueue ) > 0 {
f , hasMoreRetransmissions := s .maybeGetRetransmission (maxBytes , v )
if f != nil || hasMoreRetransmissions {
if f == nil {
return nil , nil , true
}
return f , nil , true
}
}
if len (s .dataForWriting ) == 0 && s .bufferedWriteLen () == 0 {
if s .finishedWriting && !s .finSent {
s .finSent = true
return &wire .StreamFrame {
StreamID : s .streamID ,
Offset : s .writeOffset ,
DataLenPresent : true ,
Fin : true ,
}, nil , false
}
return nil , nil , false
}
maxDataLen := s .flowController .SendWindowSize ()
if maxDataLen == 0 {
return nil , nil , true
}
reliableOffset := s .reliableOffset ()
if s .resetErr != nil && reliableOffset > 0 {
maxDataLen = min (maxDataLen , reliableOffset -s .writeOffset )
}
f , hasMoreData := s .popNewStreamFrame (maxBytes , maxDataLen , v )
if f == nil {
return nil , nil , hasMoreData
}
if f .DataLen () > 0 {
s .writeOffset += f .DataLen ()
s .flowController .AddBytesSent (f .DataLen ())
}
if s .resetErr != nil && s .writeOffset >= reliableOffset {
hasMoreData = false
}
var blocked *wire .StreamDataBlockedFrame
if f .DataLen () == maxDataLen && s .flowController .IsNewlyBlocked () {
blocked = &wire .StreamDataBlockedFrame {StreamID : s .streamID , MaximumStreamData : s .writeOffset }
}
f .Fin = s .finishedWriting && s .dataForWriting == nil && s .bufferedWriteLen () == 0 && !s .finSent
if f .Fin {
s .finSent = true
}
return f , blocked , hasMoreData
}
func (s *SendStream ) popNewStreamFrame (maxBytes , maxDataLen protocol .ByteCount , v protocol .Version ) (_ *wire .StreamFrame , hasMoreData bool ) {
f := wire .GetStreamFrame ()
f .Fin = false
f .StreamID = s .streamID
f .Offset = s .writeOffset
f .DataLenPresent = true
f .Data = f .Data [:0 ]
maxDataLen = min (maxDataLen , f .MaxDataLen (maxBytes , v ))
if maxDataLen == 0 {
f .PutBack ()
return nil , true
}
if n := min (s .bufferedWriteLen (), int (maxDataLen )); n > 0 {
f .Data = f .Data [:n ]
copy (f .Data , s .writeBuffer [s .writeBufferHead :s .writeBufferHead +n ])
s .writeBufferHead += n
if s .writeBufferHead == len (s .writeBuffer ) {
s .writeBuffer = s .writeBuffer [:0 ]
s .writeBufferHead = 0
}
s .signalWrite ()
hasMoreData = s .bufferedWriteLen () > 0 || s .dataForWriting != nil
} else {
s .getDataForWriting (f , maxDataLen )
hasMoreData = s .dataForWriting != nil || s .finishedWriting
}
if len (f .Data ) == 0 && !f .Fin {
f .PutBack ()
return nil , hasMoreData
}
return f , hasMoreData
}
func (s *SendStream ) maybeGetRetransmission (maxBytes protocol .ByteCount , v protocol .Version ) (*wire .StreamFrame , bool ) {
f := s .retransmissionQueue [0 ]
newFrame , needsSplit := f .MaybeSplitOffFrame (maxBytes , v )
if needsSplit {
return newFrame , true
}
s .retransmissionQueue = s .retransmissionQueue [1 :]
return f , len (s .retransmissionQueue ) > 0
}
func (s *SendStream ) getDataForWriting (f *wire .StreamFrame , maxBytes protocol .ByteCount ) {
if protocol .ByteCount (len (s .dataForWriting )) <= maxBytes {
f .Data = f .Data [:len (s .dataForWriting )]
copy (f .Data , s .dataForWriting )
s .dataForWriting = nil
s .signalWrite ()
return
}
f .Data = f .Data [:maxBytes ]
copy (f .Data , s .dataForWriting )
s .dataForWriting = s .dataForWriting [maxBytes :]
if s .canBufferWrite () {
s .signalWrite ()
}
}
func (s *SendStream ) isNewlyCompleted () bool {
if s .completed {
return false
}
if s .bufferedWriteLen () > 0 {
return false
}
if s .numOutstandingFrames > 0 || len (s .retransmissionQueue ) > 0 || s .queuedResetStreamFrame != nil {
return false
}
if s .finSent {
s .completed = true
return true
}
if s .resetErr != nil && (s .cancellationFlagged || s .finishedWriting ) {
s .completed = true
return true
}
return false
}
func (s *SendStream ) Close () error {
s .mutex .Lock ()
if s .shutdownErr != nil || s .finishedWriting {
s .mutex .Unlock ()
return nil
}
s .finishedWriting = true
cancelled := s .resetErr != nil
if cancelled {
s .cancellationFlagged = true
}
completed := s .isNewlyCompleted ()
s .mutex .Unlock ()
if completed {
s .sender .onStreamCompleted (s .streamID )
}
if cancelled {
return fmt .Errorf ("close called for canceled stream %d" , s .streamID )
}
s .notifyHasStreamData ()
s .ctxCancel (nil )
return nil
}
func (s *SendStream ) SetReliableBoundary () {
s .mutex .Lock ()
defer s .mutex .Unlock ()
s .reliableSize = s .writeOffset
s .reliableSize += protocol .ByteCount (s .bufferedWriteLen ())
}
func (s *SendStream ) returnFramesToPool () {
s .stopActivationTimerLocked ()
s .burstUntil = 0
s .corkPending = false
for _ , f := range s .retransmissionQueue {
f .PutBack ()
}
clear (s .retransmissionQueue )
s .retransmissionQueue = nil
s .writeBuffer = nil
s .writeBufferHead = 0
}
func (s *SendStream ) CancelWrite (errorCode StreamErrorCode ) {
s .mutex .Lock ()
if s .shutdownErr != nil {
s .mutex .Unlock ()
return
}
s .cancellationFlagged = true
if s .resetErr != nil {
completed := s .isNewlyCompleted ()
s .mutex .Unlock ()
if completed {
s .sender .onStreamCompleted (s .streamID )
}
return
}
s .resetErr = &StreamError {StreamID : s .streamID , ErrorCode : errorCode , Remote : false }
s .ctxCancel (s .resetErr )
reliableOffset := s .reliableOffset ()
if reliableOffset == 0 {
s .numOutstandingFrames = 0
s .returnFramesToPool ()
}
s .queuedResetStreamFrame = &wire .ResetStreamFrame {
StreamID : s .streamID ,
FinalSize : max (s .writeOffset , reliableOffset ),
ErrorCode : errorCode ,
ReliableSize : reliableOffset ,
}
if reliableOffset > 0 {
if buffered := protocol .ByteCount (s .bufferedWriteLen ()); buffered > 0 {
keep := reliableOffset - s .writeOffset
if keep <= 0 {
s .writeBuffer = s .writeBuffer [:0 ]
s .writeBufferHead = 0
} else if keep < buffered {
s .writeBuffer = s .writeBuffer [:s .writeBufferHead +int (keep )]
}
}
if len (s .retransmissionQueue ) > 0 {
retransmissionQueue := make ([]*wire .StreamFrame , 0 , len (s .retransmissionQueue ))
for _ , f := range s .retransmissionQueue {
if f .Offset >= reliableOffset {
f .PutBack ()
continue
}
if f .Offset +f .DataLen () <= reliableOffset {
retransmissionQueue = append (retransmissionQueue , f )
} else {
f .Data = f .Data [:reliableOffset -f .Offset ]
retransmissionQueue = append (retransmissionQueue , f )
}
}
s .retransmissionQueue = retransmissionQueue
}
}
s .mutex .Unlock ()
s .signalWrite ()
s .sender .onHasStreamControlFrame (s .streamID , s )
}
func (s *SendStream ) enableResetStreamAt () {
s .mutex .Lock ()
s .supportsResetStreamAt = true
s .mutex .Unlock ()
}
func (s *SendStream ) updateSendWindow (limit protocol .ByteCount ) {
updated := s .flowController .UpdateSendWindow (limit )
if !updated {
return
}
s .mutex .Lock ()
hasStreamData := s .dataForWriting != nil || s .bufferedWriteLen () > 0
s .mutex .Unlock ()
if hasStreamData {
s .notifyHasStreamData ()
}
}
func (s *SendStream ) onConnectionSendWindowUpdated () {
s .mutex .Lock ()
hasStreamData := s .dataForWriting != nil || s .bufferedWriteLen () > 0
s .mutex .Unlock ()
if hasStreamData {
s .notifyHasStreamData ()
}
}
func (s *SendStream ) handleStopSendingFrame (f *wire .StopSendingFrame ) {
s .mutex .Lock ()
if s .shutdownErr != nil {
s .mutex .Unlock ()
return
}
if s .resetErr != nil && s .reliableOffset () == 0 {
s .mutex .Unlock ()
return
}
s .reliableSize = 0
s .numOutstandingFrames = 0
s .returnFramesToPool ()
if s .resetErr == nil {
s .resetErr = &StreamError {StreamID : s .streamID , ErrorCode : f .ErrorCode , Remote : true }
s .ctxCancel (s .resetErr )
}
s .queuedResetStreamFrame = &wire .ResetStreamFrame {
StreamID : s .streamID ,
FinalSize : s .writeOffset ,
ErrorCode : s .resetErr .ErrorCode ,
}
s .mutex .Unlock ()
s .signalWrite ()
s .sender .onHasStreamControlFrame (s .streamID , s )
}
func (s *SendStream ) getControlFrame (monotime .Time ) (_ ackhandler .Frame , ok , hasMore bool ) {
s .mutex .Lock ()
defer s .mutex .Unlock ()
if s .queuedResetStreamFrame == nil {
return ackhandler .Frame {}, false , false
}
s .numOutstandingFrames ++
f := ackhandler .Frame {
Frame : s .queuedResetStreamFrame ,
Handler : (*sendStreamResetStreamHandler )(s ),
}
s .queuedResetStreamFrame = nil
return f , true , false
}
func (s *SendStream ) reliableOffset () protocol .ByteCount {
if !s .supportsResetStreamAt {
return 0
}
return s .reliableSize
}
func (s *SendStream ) Context () context .Context {
return s .ctx
}
func (s *SendStream ) SetWriteDeadline (t time .Time ) error {
s .mutex .Lock ()
s .deadline = monotime .FromTime (t )
s .mutex .Unlock ()
s .signalWrite ()
return nil
}
func (s *SendStream ) closeForShutdown (err error ) {
s .mutex .Lock ()
if s .shutdownErr == nil && !s .finishedWriting {
s .shutdownErr = err
s .returnFramesToPool ()
}
s .mutex .Unlock ()
s .signalWrite ()
}
func (s *SendStream ) signalWrite () {
select {
case s .writeChan <- struct {}{}:
default :
}
}
type sendStreamAckHandler SendStream
var _ ackhandler .FrameHandler = &sendStreamAckHandler {}
func (s *sendStreamAckHandler ) OnAcked (f wire .Frame ) {
sf := f .(*wire .StreamFrame )
sf .PutBack ()
s .mutex .Lock ()
if s .resetErr != nil && (*SendStream )(s ).reliableOffset () == 0 {
s .mutex .Unlock ()
return
}
s .numOutstandingFrames --
if s .numOutstandingFrames < 0 {
panic ("numOutStandingFrames negative" )
}
completed := (*SendStream )(s ).isNewlyCompleted ()
s .mutex .Unlock ()
if completed {
s .sender .onStreamCompleted (s .streamID )
}
}
func (s *sendStreamAckHandler ) OnLost (f wire .Frame ) {
sf := f .(*wire .StreamFrame )
s .mutex .Lock ()
if s .resetErr != nil && (*SendStream )(s ).reliableOffset () == 0 {
sf .PutBack ()
s .mutex .Unlock ()
return
}
s .numOutstandingFrames --
if s .numOutstandingFrames < 0 {
panic ("numOutStandingFrames negative" )
}
if s .resetErr != nil && (*SendStream )(s ).reliableOffset () > 0 {
if sf .Offset >= (*SendStream )(s ).reliableOffset () {
sf .PutBack ()
completed := (*SendStream )(s ).isNewlyCompleted ()
s .mutex .Unlock ()
if completed {
s .sender .onStreamCompleted (s .streamID )
}
return
}
if sf .Offset +sf .DataLen () > (*SendStream )(s ).reliableOffset () {
sf .Data = sf .Data [:(*SendStream )(s ).reliableOffset ()-sf .Offset ]
}
}
sf .DataLenPresent = true
s .retransmissionQueue = append (s .retransmissionQueue , sf )
s .mutex .Unlock ()
(*SendStream )(s ).notifyHasStreamData ()
}
type sendStreamResetStreamHandler SendStream
var _ ackhandler .FrameHandler = &sendStreamResetStreamHandler {}
func (s *sendStreamResetStreamHandler ) OnAcked (f wire .Frame ) {
rsf := f .(*wire .ResetStreamFrame )
s .mutex .Lock ()
if rsf .ReliableSize != (*SendStream )(s ).reliableOffset () {
s .mutex .Unlock ()
return
}
s .numOutstandingFrames --
if s .numOutstandingFrames < 0 {
panic ("numOutStandingFrames negative" )
}
completed := (*SendStream )(s ).isNewlyCompleted ()
s .mutex .Unlock ()
if completed {
s .sender .onStreamCompleted (s .streamID )
}
}
func (s *sendStreamResetStreamHandler ) OnLost (f wire .Frame ) {
rsf := f .(*wire .ResetStreamFrame )
s .mutex .Lock ()
if rsf .ReliableSize != (*SendStream )(s ).reliableOffset () {
s .mutex .Unlock ()
return
}
s .queuedResetStreamFrame = rsf
s .numOutstandingFrames --
s .mutex .Unlock ()
s .sender .onHasStreamControlFrame (s .streamID , (*SendStream )(s ))
}
The pages are generated with Golds v0.8.4 . (GOOS=linux GOARCH=amd64)
Golds is a Go 101 project developed by Tapir Liu .
PR and bug reports are welcome and can be submitted to the issue list .
Please follow @zigo_101 (reachable from the left QR code) to get the latest news of Golds .