@@ -426,7 +426,7 @@ func (r *retrier) retryWriteReadLocked(buf []byte) (int, error) {
426426
427427 // all of buf was written to c
428428 // require a response within a short timeout on r.conn (same as newConn)
429- r . shorterReadDeadlineForRetryLocked ( )
429+ newConn . SetReadDeadline ( time . Now (). Add ( r . readTimeoutLocked ()) )
430430 return newConn .Read (buf )
431431}
432432
@@ -489,12 +489,13 @@ func (r *retrier) Read(buf []byte) (n int, err error) {
489489 r .dialerID (), r .retryCount , len (r .dialers ), c , r .nextDialerIdx , laddr (c ), r .raddr , core .FmtPeriod (r .timeout ), n , len (buf ), retryReadErr )
490490 }
491491 if c != nil && core .IsNotNil (c ) {
492- // caller might have set read or write deadlines before the retry
492+ // caller might have set read or write deadlines before the retry;
493+ // if not, clear any deadlines set by the retrier
493494 _ = c .SetReadDeadline (r .readDeadline )
494495 _ = c .SetWriteDeadline (r .writeDeadline )
495496 }
496- logeor (err , note )("retrier: read: %s: #%d + (mult? %d / %d) [%s<=%s]; t : %s; b: %d/%d; err? %v" ,
497- r .dialerID (), r .retryCount , len (r .dialers ), r .nextDialerIdx , laddr (c ), r .raddr , core .FmtPeriod (r .timeout ), n , len (buf ), err )
497+ logeor (err , note )("retrier: read: %s: #%d + (mult? %d / %d) [%s<=%s]; rshortt: %s / rfullt : %s; b: %d/%d; err? %v" ,
498+ r .dialerID (), r .retryCount , len (r .dialers ), r .nextDialerIdx , laddr (c ), r .raddr , core .FmtPeriod (r .timeout ), core . FmtTimeAsPeriod ( r . readDeadline ), n , len (buf ), err )
498499 r .tee = nil // discard teed data
499500 return
500501 }
@@ -504,7 +505,7 @@ func (r *retrier) Read(buf []byte) (n int, err error) {
504505 return
505506}
506507
507- func (r * retrier ) teedFirstWrite (b []byte ) (n int , firstWrite , didAttemptWrite bool , src net.Addr , err error ) {
508+ func (r * retrier ) teedFirstWrite (b []byte ) (n int , firstWrite , didAttemptWrite bool , readWait time. Duration , src net.Addr , err error ) {
508509 r .mu .Lock ()
509510 defer r .mu .Unlock ()
510511
@@ -519,33 +520,24 @@ func (r *retrier) teedFirstWrite(b []byte) (n int, firstWrite, didAttemptWrite b
519520 }
520521
521522 src = laddr (c )
522- if ! r .retryCompleted () { // first write
523+ if ! r .retryCompleted () { // may be first write
523524 _ = c .SetWriteDeadline (r .writeDeadline )
524525
525526 n , err = c .Write (b )
526527
527528 // capture first write, aka "hello"
528529 r .tee = append (r .tee , b ... )
529530 didAttemptWrite = true
531+ readWait = r .readTimeoutLocked ()
530532 // all of b was written to r.tee if not to c
531533 // require a response or another write within a short timeout.
532- r . shorterReadDeadlineForRetryLocked ( )
534+ c . SetReadDeadline ( time . Now (). Add ( readWait ) )
533535 }
534536
535537 return
536538}
537539
538- func (r * retrier ) shorterReadDeadlineForRetryLocked () {
539- c := r .conn
540- if r .timeout > 0 {
541- _ = c .SetReadDeadline (time .Now ().Add (r .timeout ))
542- } else {
543- // if timeout is set to 0, then use client requested deadline
544- _ = c .SetReadDeadline (r .readDeadline )
545- }
546- }
547-
548- func (r * retrier ) readTimeout () time.Duration {
540+ func (r * retrier ) readTimeoutLocked () time.Duration {
549541 if r .timeout > 0 {
550542 return r .timeout
551543 }
@@ -563,7 +555,7 @@ func (r *retrier) Write(b []byte) (int, error) {
563555 // empty at steady-state.
564556 if ! r .retryCompleted () {
565557 // todo: what if sentAndCopied is false and err != nil?
566- n , first , sentAndCopied , src , err := r .teedFirstWrite (b )
558+ n , first , sentAndCopied , until , src , err := r .teedFirstWrite (b )
567559
568560 note := log .D
569561 if sentAndCopied {
@@ -587,33 +579,31 @@ func (r *retrier) Write(b []byte) (int, error) {
587579 // by the retry procedure. Block until we have a final socket (which will
588580 // already have replayed r.tee), and retry.
589581 // ie, wait until first write is done on the final socket.
590- until := r .readTimeout ()
591582 maxUntil := max (until , until * maxRetryCount )
592583 if r .multidial {
593584 maxUntil = max (maxUntil , maxUntil * time .Duration (len (r .dialers )))
594585 }
595586 select {
596587 case <- r .retryDoneCh :
597588 case <- time .After (maxUntil ): // arb high timeout; it should rarely if ever needed
598- log .W ("retrier: write: %s: 1st write timed-out waiting for %s [calc-rtt: %s] 1st read b/w [%s=>%s], mult: %d, b: %d/%d, err: %v" ,
589+ rerr := log .EE ("retrier: write: %s: 1st write timed-out waiting for %s [calc-rtt: %s] 1st read b/w [%s=>%s], mult: %d, b: %d/%d, err: %v" ,
599590 r .dialerID (), core .FmtPeriod (maxUntil ), core .FmtPeriod (r .timeout ), src , r .raddr , len (r .dialers ), n , len (b ), err )
600- return n , core .JoinErr (err , errRetryTimeout )
591+ return n , core .JoinErr (err , rerr , errRetryTimeout )
601592 }
602593
603594 r .mu .Lock ()
604595 defer r .mu .Unlock ()
605596
606- elapsed := time .Since (start ).Milliseconds ()
607597 // r.conn may be nil or closed by the time we get here
608598 finalConn := r .conn
609599 noconn := finalConn == nil || core .IsNil (finalConn )
610600 if r .retryWriteErr != nil || noconn { // check if retried writes also failed
611601 if noconn {
612602 err = core .JoinErr (err , errNilConn )
613603 }
614- log .E ("retrier: write: %s: retry failed [%s=>%s] b: %d/%d (tee: %d) in %dms ; old => new: %v => %v; noconn? %t" ,
615- r .dialerID (), laddr (r .conn ), r .raddr , n , len (b ), len (r .tee ), elapsed , err , r .retryWriteErr , noconn )
616- return n , core .JoinErr (err , r .retryWriteErr ) // pass on the og error, too
604+ werr := log .EE ("retrier: write: %s: retry failed [%s=>%s] b: %d/%d (tee: %d) in %s ; old => new: %v => %v; noconn? %t" ,
605+ r .dialerID (), laddr (r .conn ), r .raddr , n , len (b ), len (r .tee ), core . FmtTimeAsPeriod ( start ) , err , r .retryWriteErr , noconn )
606+ return n , core .JoinErr (err , r .retryWriteErr , werr ) // pass on the og error, too
617607 }
618608
619609 // if len(leftover) > 0 {
@@ -629,9 +619,9 @@ func (r *retrier) Write(b []byte) (int, error) {
629619
630620 // retryCompleted() is true, so r.conn is final and doesn't need locking
631621 if c := r .conn ; c == nil || core .IsNil (c ) {
632- log .E ("retrier: write: %s: [] => %s (b: %d, tee: %d), not retrying, but no conn" ,
622+ cerr := log .EE ("retrier: write: %s: [] => %s (b: %d, tee: %d), not retrying, but no conn" ,
633623 r .dialerID (), r .raddr , len (b ), len (r .tee ))
634- return 0 , errNilConn
624+ return 0 , core . JoinErr ( cerr , errNilConn )
635625 } else {
636626 return c .Write (b )
637627 }
@@ -672,16 +662,22 @@ func (r *retrier) ReadFrom(reader io.Reader) (bytes int64, err error) {
672662 }
673663 }
674664
665+ // disable read and write deadlines as io.ReaderFrom does not
666+ // rely on io.Read and io.Write semantics for "r.conn" from which
667+ // deadlines are extended to avoid timeouts (see also: rwconn.go)
675668 optimizedReadFrom := true
676669 var b int64
677670 switch x := c .(type ) {
678671 case * net.TCPConn :
672+ r .SetDeadline (time.Time {})
679673 b , err = x .ReadFrom (reader )
680674 bytes += b
681675 case * splitter :
676+ r .SetDeadline (time.Time {})
682677 b , err = x .ReadFrom (reader )
683678 bytes += b
684679 case io.ReaderFrom :
680+ r .SetDeadline (time.Time {})
685681 b , err = x .ReadFrom (reader )
686682 bytes += b
687683 default : // net.UDPConn, net.PacketConn etc?
@@ -691,13 +687,6 @@ func (r *retrier) ReadFrom(reader io.Reader) (bytes int64, err error) {
691687 bytes += b
692688 }
693689
694- if optimizedReadFrom {
695- // disable read and write deadlines as io.ReaderFrom does not
696- // rely on io.Read and io.Write semantics from which deadlines
697- // are usually extended to avoid timeouts (see also: rwconn.go)
698- r .SetDeadline (time.Time {})
699- }
700-
701690 logeif (err )("retrier: readfrom: %s: (optimized? %t for %T) done (id: %s, pinned? %t); sz: %d; err: %v" ,
702691 r .dialerID (), optimizedReadFrom , c , pinnedID , pinned , bytes , err )
703692 return
0 commit comments