From 4d89d732e211dfa9b584215be3d1d413aa0da48b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=B8=96=E7=95=8C?= Date: Fri, 13 Mar 2026 20:07:18 +0800 Subject: [PATCH] ccm,ocm: reset connector backoff after successful connection The consecutiveFailures counter in connectorLoop never resets, causing backoff to permanently cap at 30-45s even after a connection that served successfully for hours. Reset the counter when connectorConnect ran for at least one minute, indicating a successful session rather than a transient dial/handshake failure. --- service/ccm/reverse.go | 31 +++++++++++++++++++------------ service/ocm/reverse.go | 31 +++++++++++++++++++------------ 2 files changed, 38 insertions(+), 24 deletions(-) diff --git a/service/ccm/reverse.go b/service/ccm/reverse.go index 7e38c9ced..97d71f0ad 100644 --- a/service/ccm/reverse.go +++ b/service/ccm/reverse.go @@ -132,10 +132,13 @@ func (c *externalCredential) connectorLoop() { default: } - err := c.connectorConnect() + sessionLifetime, err := c.connectorConnect() if c.reverseContext.Err() != nil { return } + if sessionLifetime >= connectorBackoffResetThreshold { + consecutiveFailures = 0 + } consecutiveFailures++ backoff := connectorBackoff(consecutiveFailures) c.logger.Warn("reverse connection for ", c.tag, " lost: ", err, ", reconnecting in ", backoff) @@ -147,6 +150,8 @@ func (c *externalCredential) connectorLoop() { } } +const connectorBackoffResetThreshold = time.Minute + func connectorBackoff(failures int) time.Duration { if failures > 5 { failures = 5 @@ -159,14 +164,14 @@ func connectorBackoff(failures int) time.Duration { return base + jitter } -func (c *externalCredential) connectorConnect() error { +func (c *externalCredential) connectorConnect() (time.Duration, error) { if c.reverseService == nil { - return E.New("reverse service not initialized") + return 0, E.New("reverse service not initialized") } destination := c.connectorResolveDestination() conn, err := c.connectorDialer.DialContext(c.reverseContext, "tcp", destination) if err != nil { - return E.Cause(err, "dial") + return 0, E.Cause(err, "dial") } if c.connectorTLS != nil { @@ -174,7 +179,7 @@ func (c *externalCredential) connectorConnect() error { err = tlsConn.HandshakeContext(c.reverseContext) if err != nil { conn.Close() - return E.Cause(err, "tls handshake") + return 0, E.Cause(err, "tls handshake") } conn = tlsConn } @@ -188,24 +193,24 @@ func (c *externalCredential) connectorConnect() error { _, err = io.WriteString(conn, upgradeRequest) if err != nil { conn.Close() - return E.Cause(err, "write upgrade request") + return 0, E.Cause(err, "write upgrade request") } reader := bufio.NewReader(conn) statusLine, err := reader.ReadString('\n') if err != nil { conn.Close() - return E.Cause(err, "read upgrade response") + return 0, E.Cause(err, "read upgrade response") } if !strings.HasPrefix(statusLine, "HTTP/1.1 101") { conn.Close() - return E.New("unexpected upgrade response: ", strings.TrimSpace(statusLine)) + return 0, E.New("unexpected upgrade response: ", strings.TrimSpace(statusLine)) } for { line, readErr := reader.ReadString('\n') if readErr != nil { conn.Close() - return E.Cause(readErr, "read upgrade headers") + return 0, E.Cause(readErr, "read upgrade headers") } if strings.TrimSpace(line) == "" { break @@ -215,22 +220,24 @@ func (c *externalCredential) connectorConnect() error { session, err := yamux.Server(conn, reverseYamuxConfig()) if err != nil { conn.Close() - return E.Cause(err, "create yamux server") + return 0, E.Cause(err, "create yamux server") } defer session.Close() c.logger.Info("reverse connection established for ", c.tag) + serveStart := time.Now() httpServer := &http.Server{ Handler: c.reverseService, ReadTimeout: 0, IdleTimeout: 120 * time.Second, } err = httpServer.Serve(&yamuxNetListener{session: session}) + sessionLifetime := time.Since(serveStart) if err != nil && !errors.Is(err, http.ErrServerClosed) && c.reverseContext.Err() == nil { - return E.Cause(err, "serve") + return sessionLifetime, E.Cause(err, "serve") } - return E.New("connection closed") + return sessionLifetime, E.New("connection closed") } func (c *externalCredential) connectorResolveDestination() M.Socksaddr { diff --git a/service/ocm/reverse.go b/service/ocm/reverse.go index 23ca1cc47..b3c17f45f 100644 --- a/service/ocm/reverse.go +++ b/service/ocm/reverse.go @@ -132,10 +132,13 @@ func (c *externalCredential) connectorLoop() { default: } - err := c.connectorConnect() + sessionLifetime, err := c.connectorConnect() if c.reverseContext.Err() != nil { return } + if sessionLifetime >= connectorBackoffResetThreshold { + consecutiveFailures = 0 + } consecutiveFailures++ backoff := connectorBackoff(consecutiveFailures) c.logger.Warn("reverse connection for ", c.tag, " lost: ", err, ", reconnecting in ", backoff) @@ -147,6 +150,8 @@ func (c *externalCredential) connectorLoop() { } } +const connectorBackoffResetThreshold = time.Minute + func connectorBackoff(failures int) time.Duration { if failures > 5 { failures = 5 @@ -159,14 +164,14 @@ func connectorBackoff(failures int) time.Duration { return base + jitter } -func (c *externalCredential) connectorConnect() error { +func (c *externalCredential) connectorConnect() (time.Duration, error) { if c.reverseService == nil { - return E.New("reverse service not initialized") + return 0, E.New("reverse service not initialized") } destination := c.connectorResolveDestination() conn, err := c.connectorDialer.DialContext(c.reverseContext, "tcp", destination) if err != nil { - return E.Cause(err, "dial") + return 0, E.Cause(err, "dial") } if c.connectorTLS != nil { @@ -174,7 +179,7 @@ func (c *externalCredential) connectorConnect() error { err = tlsConn.HandshakeContext(c.reverseContext) if err != nil { conn.Close() - return E.Cause(err, "tls handshake") + return 0, E.Cause(err, "tls handshake") } conn = tlsConn } @@ -188,24 +193,24 @@ func (c *externalCredential) connectorConnect() error { _, err = io.WriteString(conn, upgradeRequest) if err != nil { conn.Close() - return E.Cause(err, "write upgrade request") + return 0, E.Cause(err, "write upgrade request") } reader := bufio.NewReader(conn) statusLine, err := reader.ReadString('\n') if err != nil { conn.Close() - return E.Cause(err, "read upgrade response") + return 0, E.Cause(err, "read upgrade response") } if !strings.HasPrefix(statusLine, "HTTP/1.1 101") { conn.Close() - return E.New("unexpected upgrade response: ", strings.TrimSpace(statusLine)) + return 0, E.New("unexpected upgrade response: ", strings.TrimSpace(statusLine)) } for { line, readErr := reader.ReadString('\n') if readErr != nil { conn.Close() - return E.Cause(readErr, "read upgrade headers") + return 0, E.Cause(readErr, "read upgrade headers") } if strings.TrimSpace(line) == "" { break @@ -215,22 +220,24 @@ func (c *externalCredential) connectorConnect() error { session, err := yamux.Server(conn, reverseYamuxConfig()) if err != nil { conn.Close() - return E.Cause(err, "create yamux server") + return 0, E.Cause(err, "create yamux server") } defer session.Close() c.logger.Info("reverse connection established for ", c.tag) + serveStart := time.Now() httpServer := &http.Server{ Handler: c.reverseService, ReadTimeout: 0, IdleTimeout: 120 * time.Second, } err = httpServer.Serve(&yamuxNetListener{session: session}) + sessionLifetime := time.Since(serveStart) if err != nil && !errors.Is(err, http.ErrServerClosed) && c.reverseContext.Err() == nil { - return E.Cause(err, "serve") + return sessionLifetime, E.Cause(err, "serve") } - return E.New("connection closed") + return sessionLifetime, E.New("connection closed") } func (c *externalCredential) connectorResolveDestination() M.Socksaddr {