Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions cmd/replayer/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"github.com/pingcap/tiproxy/lib/config"
"github.com/pingcap/tiproxy/lib/util/cmd"
"github.com/pingcap/tiproxy/pkg/manager/cert"
"github.com/pingcap/tiproxy/pkg/manager/health"
"github.com/pingcap/tiproxy/pkg/manager/id"
"github.com/pingcap/tiproxy/pkg/manager/logger"
"github.com/pingcap/tiproxy/pkg/manager/memory"
Expand Down Expand Up @@ -85,6 +86,17 @@ func main() {
cfgMgr := &nopConfigManager{cfg: cfg}
memMgr := memory.NewMemManager(lg, cfgMgr)
memMgr.Start(context.Background())
healthMgr := health.NewManager(
func() bool { return true },
func() (bool, string) {
reject, snapshot, threshold := memMgr.ShouldRejectNewConn()
if !reject {
return false, ""
}
return true, fmt.Sprintf("high memory usage (usage=%.4f, threshold=%.4f, used=%d, limit=%d, last_update=%s)",
snapshot.Usage, threshold, snapshot.Used, snapshot.Limit, snapshot.UpdateTime.String())
},
)

// create replay job manager
hsHandler := backend.NewStaticHandshakeHandler(*addr)
Expand All @@ -98,6 +110,7 @@ func main() {
CertMgr: cert.NewCertManager(),
BackendReader: nil,
ReplayJobMgr: r,
Health: healthMgr,
}
var ready atomic.Bool
ready.Store(true)
Expand Down
65 changes: 65 additions & 0 deletions pkg/manager/health/health.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
// Copyright 2026 PingCAP, Inc.
// SPDX-License-Identifier: Apache-2.0

package health

import "sync/atomic"

// Manager aggregates the signals that decide whether the TiProxy instance is
// serving and whether it should keep accepting new connections. Each public
// method below evaluates its own condition independently, so reordering one
// signal never silently changes another consumer's behavior.
type Manager struct {
// shuttingDown means the instance is in graceful shutdown. It only affects
// Healthy (so DebugHealth reports unhealthy and the LB drains); it does NOT
// reject new connections, since the proxy keeps serving until its listeners
// are closed.
shuttingDown atomic.Bool
ready func() bool
rejectCheck func() (bool, string)
}

// NewManager creates a Manager.
// - ready reports whether the namespace manager is ready (init phase when false).
// - rejectCheck reports whether new connections should be rejected because of
// memory pressure, together with a human-readable reason.
func NewManager(ready func() bool, rejectCheck func() (bool, string)) *Manager {
return &Manager{ready: ready, rejectCheck: rejectCheck}
}

// PreClose marks the instance as gracefully shutting down. Idempotent; safe to
// call from PreClose paths.
func (m *Manager) PreClose() {
m.shuttingDown.Store(true)
}

// Healthy reports whether the instance is fully serving. Used by DebugHealth.
// It returns false (with a reason) during init, rejectConns, and graceful
// shutdown.
func (m *Manager) Healthy() (bool, string) {
if m.shuttingDown.Load() {
return false, "server is shutting down"
}
if m.rejectCheck != nil {
if reject, reason := m.rejectCheck(); reject {
return false, reason
}
}
if m.ready != nil && !m.ready() {
return false, "server is not ready"
}
return true, ""
}

// RejectConns reports whether new connections should be rejected and returns a
// reason string. Used by the proxy server. It returns true only on memory
// pressure; graceful shutdown does NOT reject here, because the proxy keeps
// accepting until its listeners are closed.
func (m *Manager) RejectConns() (bool, string) {
if m.rejectCheck != nil {
if reject, reason := m.rejectCheck(); reject {
return true, reason
}
}
return false, ""
}
79 changes: 79 additions & 0 deletions pkg/manager/health/health_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
// Copyright 2026 PingCAP, Inc.
// SPDX-License-Identifier: Apache-2.0

package health

import (
"testing"

"github.com/stretchr/testify/require"
)

func TestManager(t *testing.T) {
ready := true
reject := false
mgr := NewManager(
func() bool { return ready },
func() (bool, string) {
if !reject {
return false, ""
}
return true, "high memory usage"
},
)

// Fully serving.
ok, reason := mgr.Healthy()
require.True(t, ok)
require.Empty(t, reason)
rejectConn, reason := mgr.RejectConns()
require.False(t, rejectConn)
require.Empty(t, reason)

// Init phase: Healthy reports not ready, but the proxy still accepts.
ready = false
ok, reason = mgr.Healthy()
require.False(t, ok)
require.Equal(t, "server is not ready", reason)
rejectConn, reason = mgr.RejectConns()
require.False(t, rejectConn)
require.Empty(t, reason)

// Memory pressure: both consumers see the reject reason.
ready = true
reject = true
ok, reason = mgr.Healthy()
require.False(t, ok)
require.Equal(t, "high memory usage", reason)
rejectConn, reason = mgr.RejectConns()
require.True(t, rejectConn)
require.Equal(t, "high memory usage", reason)

// Graceful shutdown alone (no memory pressure): Healthy reports unhealthy,
// but the proxy keeps accepting until its listeners are closed.
reject = false
mgr.PreClose()
ok, reason = mgr.Healthy()
require.False(t, ok)
require.Equal(t, "server is shutting down", reason)
rejectConn, reason = mgr.RejectConns()
require.False(t, rejectConn)
require.Empty(t, reason)
}

func TestManagerNilChecks(t *testing.T) {
// Nil ready/rejectCheck must not panic and default to serving/accepting.
mgr := NewManager(nil, nil)
ok, reason := mgr.Healthy()
require.True(t, ok)
require.Empty(t, reason)
rejectConn, reason := mgr.RejectConns()
require.False(t, rejectConn)
require.Empty(t, reason)

mgr.PreClose()
ok, _ = mgr.Healthy()
require.False(t, ok)
rejectConn, _ = mgr.RejectConns()
require.False(t, rejectConn)
}
81 changes: 39 additions & 42 deletions pkg/proxy/proxy.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@ import (
"github.com/pingcap/tiproxy/pkg/balance/router"
"github.com/pingcap/tiproxy/pkg/manager/cert"
"github.com/pingcap/tiproxy/pkg/manager/id"
mgrmem "github.com/pingcap/tiproxy/pkg/manager/memory"
"github.com/pingcap/tiproxy/pkg/metrics"
"github.com/pingcap/tiproxy/pkg/proxy/backend"
"github.com/pingcap/tiproxy/pkg/proxy/client"
Expand Down Expand Up @@ -47,24 +46,28 @@ type BackendDialer interface {
}

type SQLServer struct {
listeners []net.Listener
addrs []string
logger *zap.Logger
certMgr *cert.CertManager
idMgr *id.IDManager
memUsage memoryStateProvider
hsHandler backend.HandshakeHandler
cpt capture.Capture
meter backend.Meter
dialer BackendDialer
wg waitgroup.WaitGroup
cancelFunc context.CancelFunc
listeners []net.Listener
addrs []string
logger *zap.Logger
certMgr *cert.CertManager
idMgr *id.IDManager
connBufferUpdater connBufferMemoryUpdater
health connAcceptor
hsHandler backend.HandshakeHandler
cpt capture.Capture
meter backend.Meter
dialer BackendDialer
wg waitgroup.WaitGroup
cancelFunc context.CancelFunc

mu serverState
}

type memoryStateProvider interface {
ShouldRejectNewConn() (bool, mgrmem.UsageSnapshot, float64)
// connAcceptor reports whether the proxy should reject new connections and
// provides a reason for the reject log. Satisfied by *health.Manager; defined
// here so the proxy package does not depend on the health package.
type connAcceptor interface {
RejectConns() (bool, string)
}

type connBufferMemoryUpdater interface {
Expand All @@ -81,16 +84,17 @@ func estimateConnBufferMemDelta(bufferSize int) int64 {

// NewSQLServer creates a new SQLServer.
func NewSQLServer(logger *zap.Logger, cfg *config.Config, certMgr *cert.CertManager, idMgr *id.IDManager, cpt capture.Capture,
meter backend.Meter, hsHandler backend.HandshakeHandler, memUsage memoryStateProvider) (*SQLServer, error) {
meter backend.Meter, hsHandler backend.HandshakeHandler, connBufferUpdater connBufferMemoryUpdater, health connAcceptor) (*SQLServer, error) {
var err error
s := &SQLServer{
logger: logger,
certMgr: certMgr,
idMgr: idMgr,
memUsage: memUsage,
hsHandler: hsHandler,
cpt: cpt,
meter: meter,
logger: logger,
certMgr: certMgr,
idMgr: idMgr,
connBufferUpdater: connBufferUpdater,
health: health,
hsHandler: hsHandler,
cpt: cpt,
meter: meter,
mu: serverState{
clients: make(map[uint64]*client.ClientConnection),
},
Expand Down Expand Up @@ -182,17 +186,15 @@ func (s *SQLServer) Run(ctx context.Context, cfgch <-chan *config.Config) {
}

func (s *SQLServer) onConn(ctx context.Context, conn net.Conn, addr string) {
if s.rejectConnByMemory(conn) {
if s.rejectConn(conn) {
return
}

var (
connBufferUpdater connBufferMemoryUpdater
connBufferMemDelta int64
)
if s.memUsage != nil {
connBufferUpdater, _ = s.memUsage.(connBufferMemoryUpdater)
}
connBufferUpdater = s.connBufferUpdater

tcpKeepAlive, logger, connID, clientConn := func() (bool, *zap.Logger, uint64, *client.ClientConnection) {
s.mu.Lock()
Expand Down Expand Up @@ -265,24 +267,19 @@ func (s *SQLServer) onConn(ctx context.Context, conn net.Conn, addr string) {
clientConn.Run(ctx)
}

func (s *SQLServer) rejectConnByMemory(conn net.Conn) bool {
if s.memUsage == nil {
func (s *SQLServer) rejectConn(conn net.Conn) bool {
if s.health == nil {
return false
}
reject, snapshot, threshold := s.memUsage.ShouldRejectNewConn()
if !reject {
return false
if reject, reason := s.health.RejectConns(); reject {
metrics.RejectConnCounter.WithLabelValues("memory").Inc()
s.logger.Warn("reject connection",
zap.String("reason", reason),
zap.Stringer("client_addr", conn.RemoteAddr()),
zap.Error(conn.Close()))
return true
}
metrics.RejectConnCounter.WithLabelValues("memory").Inc()
s.logger.Warn("reject connection due to high memory usage",
zap.Stringer("client_addr", conn.RemoteAddr()),
zap.Float64("threshold", threshold),
zap.Float64("usage", snapshot.Usage),
zap.Uint64("used", snapshot.Used),
zap.Uint64("limit", snapshot.Limit),
zap.Time("last_update", snapshot.UpdateTime),
zap.Error(conn.Close()))
return true
return false
}

func (s *SQLServer) fromPublicEndpoint(addr net.Addr) bool {
Expand Down
Loading