Documentation
¶
Overview ¶
Package gproxy implements a gnet-based event-driven MTProxy server.
Index ¶
- Variables
- func CheckFrameSize(size int) bool
- func IsUnixSocket(addr string) bool
- func MiddleEndResponseProcessingBytes() int
- func Run(cfg *Config, logger Logger) (shutdown func(), errCh <-chan error)
- func Zeroize(b []byte)
- func ZeroizeArray32(arr *[32]byte)
- type BufferPool
- type Config
- type ConnContext
- func (c *ConnContext) Cleanup()
- func (c *ConnContext) DCID() int
- func (c *ConnContext) ID() uint64
- func (c *ConnContext) LogPrefix() string
- func (c *ConnContext) ProtocolMode() ProtocolMode
- func (c *ConnContext) RealClientAddr(fallback net.Addr) net.Addr
- func (c *ConnContext) Relay() *RelayContext
- func (c *ConnContext) SetProtocolMode(mode ProtocolMode)
- func (c *ConnContext) SetRealClientAddr(addr net.Addr)
- func (c *ConnContext) SetRelay(r *RelayContext)
- func (c *ConnContext) SetSpliceOverride(addr string)
- func (c *ConnContext) SetState(state ConnState)
- func (c *ConnContext) SetTrafficCounters(bytesIn, bytesOut *atomic.Int64)
- func (c *ConnContext) State() ConnState
- func (c *ConnContext) TrafficIn() *atomic.Int64
- func (c *ConnContext) TrafficOut() *atomic.Int64
- type ConnLimiter
- type ConnState
- type DCConnContext
- type DesyncDetector
- type HandshakeFailureStat
- type HotReloadConfig
- type HotReloader
- type InternalProxyAuth
- type Logger
- type LogicalQueueBudget
- type LogicalStream
- func (s *LogicalStream) Close() error
- func (s *LogicalStream) Discard(n int) (int, error)
- func (s *LogicalStream) InboundBuffered() int
- func (s *LogicalStream) LocalAddr() net.Addr
- func (s *LogicalStream) OutboundBuffered() int
- func (s *LogicalStream) OwnerStopped()
- func (s *LogicalStream) Peek(n int) ([]byte, error)
- func (s *LogicalStream) ReadableBytes() int
- func (s *LogicalStream) RemoteAddr() net.Addr
- func (s *LogicalStream) TryRead(dst []byte) (int, error)
- func (s *LogicalStream) TryWrite(data []byte) (int, error)
- func (s *LogicalStream) Write(data []byte) (int, error)
- type LogicalStreamOptions
- type MiddleEndBindingSource
- type MiddleEndFrontendConfig
- type MiddleEndFrontendStats
- type MiddleEndPrecommitAction
- type ProtocolMode
- type ProxyHandler
- func NewProxyHandler(cfg *Config, logger Logger) *ProxyHandler
- func NewProxyHandlerWithMiddleEnd(cfg *Config, logger Logger, frontend MiddleEndFrontendConfig) (*ProxyHandler, error)
- func RunWithHandler(cfg *Config, logger Logger) (shutdown func(), handler *ProxyHandler, errCh <-chan error)
- func RunWithMiddleEnd(cfg *Config, logger Logger, frontend MiddleEndFrontendConfig) (shutdown func(), handler *ProxyHandler, errCh <-chan error, err error)
- func (h *ProxyHandler) ApplyHotConfig(cfg *Config)
- func (h *ProxyHandler) HandshakeFailureStats() []HandshakeFailureStat
- func (h *ProxyHandler) IdleTimeout() time.Duration
- func (h *ProxyHandler) MiddleEndFrontendStats() MiddleEndFrontendStats
- func (h *ProxyHandler) OnBoot(eng gnet.Engine) gnet.Action
- func (h *ProxyHandler) OnClose(c gnet.Conn, err error) gnet.Action
- func (h *ProxyHandler) OnOpen(c gnet.Conn) ([]byte, gnet.Action)
- func (h *ProxyHandler) OnShutdown(eng gnet.Engine)
- func (h *ProxyHandler) OnTick() (time.Duration, gnet.Action)
- func (h *ProxyHandler) OnTraffic(c gnet.Conn) gnet.Action
- func (h *ProxyHandler) OpenLogicalStream(options LogicalStreamOptions) (*LogicalStream, error)
- func (h *ProxyHandler) UserLimiter() *UserIPLimiter
- type ProxyProtoResult
- type RelayContext
- type ReplayCache
- type Secret
- type UserIPLimiter
- func (l *UserIPLimiter) Close()
- func (l *UserIPLimiter) LimitingEnabled() bool
- func (l *UserIPLimiter) Release(key string)
- func (l *UserIPLimiter) Stats() []UserIPStats
- func (l *UserIPLimiter) TrafficCounters(secret []byte) (bytesIn, bytesOut *atomic.Int64)
- func (l *UserIPLimiter) TryAcquire(ip net.IP, secret []byte, secretName string) (key string, ok bool)
- type UserIPStats
Constants ¶
This section is empty.
Variables ¶
var ( // ErrInvalidMiddleEndFrontend reports an invalid explicit frontend // dependency or bound. ErrInvalidMiddleEndFrontend = errors.New("invalid Middle-End frontend") // ErrMiddleEndClientBackpressure reports an authenticated client whose // unread input exceeded the configured owner-loop bound. ErrMiddleEndClientBackpressure = errors.New("Middle-End client input limit exceeded") // ErrMiddleEndClientProtocol reports invalid client framing or a source // event that cannot be represented on the client stream. ErrMiddleEndClientProtocol = errors.New("Middle-End client protocol failure") )
var ErrLogicalStreamConfig = errors.New("invalid logical MTProxy stream configuration")
Functions ¶
func CheckFrameSize ¶ added in v0.1.8
CheckFrameSize checks if a frame size indicates desync. Returns true if the frame size is abnormally large (likely desync).
func IsUnixSocket ¶ added in v0.1.4
IsUnixSocket returns true if the bind address is a Unix socket.
func MiddleEndResponseProcessingBytes ¶ added in v0.6.5
func MiddleEndResponseProcessingBytes() int
MiddleEndResponseProcessingBytes is the minimum reserve for one complete frontend encode. The ME decoder has its own bounded scratch storage.
func Run ¶
Run starts the proxy with graceful shutdown support using gnet. Returns a shutdown function that can be called to stop the server.
func Zeroize ¶ added in v0.1.8
func Zeroize(b []byte)
Zeroize overwrites a byte slice with zeros. This is used for secure cleanup of sensitive data like session IDs.
Note: Go's compiler may optimize away zeroing of "dead" variables. For maximum security, call this before the variable goes out of scope while it's still reachable.
func ZeroizeArray32 ¶ added in v0.1.8
func ZeroizeArray32(arr *[32]byte)
ZeroizeArray32 zeros a 32-byte array in place.
Types ¶
type BufferPool ¶ added in v0.1.8
type BufferPool struct {
// contains filtered or unexported fields
}
BufferPool is a sync.Pool wrapper for reusable byte buffers.
func NewBufferPool ¶ added in v0.1.8
func NewBufferPool(size int) *BufferPool
NewBufferPool creates a new buffer pool with the given buffer size.
func (*BufferPool) Get ¶ added in v0.1.8
func (bp *BufferPool) Get() *[]byte
Get retrieves a buffer from the pool. The returned buffer may contain stale data - caller should slice or overwrite.
func (*BufferPool) Put ¶ added in v0.1.8
func (bp *BufferPool) Put(buf *[]byte)
Put returns a buffer to the pool. The buffer should not be used after calling Put.
type Config ¶
type Config struct {
// Secrets is the list of allowed proxy secrets.
Secrets []Secret
Host string // Default SNI hostname (from first secret)
// Network
BindAddr string
// TLS Fronting
MaskHost string // Domain to mimic (SNI validation, proxy links)
MaskPort int // Default port
FetchRealCert bool
SpliceUnrecognized bool
CertRefreshHours int
// Certificate fetching (where to connect to get real cert)
// Defaults to MaskHost:MaskPort if not set
CertHost string
CertPort int
// FakeCertSize sets the exact payload size of the fake encrypted-certificate
// ApplicationData record in the FakeTLS ServerHello. 0 = auto (match the mask
// backend's real first cert-record size, falling back to random padding).
FakeCertSize int
// MaskSNISafelist is an opt-in list of extra domains that an unauthenticated
// probe may be fronted to. When a bad-secret ClientHello's SNI is on this
// list, the connection is spliced to that domain (port 443) instead of the
// default splice target, so the on-wire conversation matches the claimed SNI.
MaskSNISafelist []string
// ClockSyncURL, when set, is fetched once at startup; its HTTP Date header
// corrects a skewed server clock for the handshake time-skew check.
ClockSyncURL string
// Splice target (where to forward unrecognized clients)
// Defaults to MaskHost:MaskPort if not set
SpliceHost string
SplicePort int
SpliceProxyProtocol int // 0 = off, 1 = v1 (text), 2 = v2 (binary)
SpliceIdleTimeout time.Duration // Idle timeout for splice connections (default 30s)
// Performance
IPPreference dc.IPPreference
IdleTimeout time.Duration
TimeSkewTolerance time.Duration
// ClientSilenceClose, when > 0, closes a relaying connection whose last
// relayed payload was server->client and has gone unanswered by the client
// for this long. Breaks an iOS MtProtoKit bad_salt "Updating" wedge. 0 = off.
ClientSilenceClose time.Duration
// Upstream (DC connection)
Socks5Addr string // SOCKS5 proxy for DC connections (e.g., "127.0.0.1:1080")
// Incoming connection handling
ProxyProtocol bool // Accept incoming PROXY protocol headers
// InternalProxyProtocol accepts PROXY headers only from loopback TCP or Unix
// peers that authenticate with the process-local token. Native WEB backends
// use it to retain the validated browser IP.
InternalProxyProtocol bool
InternalProxyAuth *InternalProxyAuth
HandshakeTimeout time.Duration // Max time for handshake before closing (default 30s)
MaxConnections int // Max concurrent accepted connections, 0 = unlimited
MaxConnectionsPerIP int // Max concurrent connections per IP+secret, 0 = unlimited
MaxIPsPerUser int // Max unique IPs per user, 0 = unlimited
IPBlockTimeout time.Duration // How long blocked IPs stay blocked (default 5m)
// Backpressure
MaxWriteBuffer int // Max pending bytes per connection before closing (0 = 4MB default)
// Anti-DPI record shaping (proxy -> client direction).
// EnableDRS turns on the Chrome-style probe-then-ramp record sizer:
// outbound TLS ApplicationData starts with 1369-byte records and ramps
// to full size after 8 records or 128 KB.
EnableDRS bool
// EnableSplitTLS emits the first outbound ApplicationData record as a
// 1-byte record before continuing normally. Defeats passive signatures
// keyed on the first record.
EnableSplitTLS bool
// gnet-specific
Multicore bool // Use multiple event loops
ReusePort bool // Enable SO_REUSEPORT
LockOSThread bool // Lock goroutines to OS threads
NumEventLoop int // Number of event loops (0 = auto)
// WebProxyFingerprint changes when restart-only [web-proxy] settings change.
// It contains no secret values and exists only for hot-reload diagnostics.
WebProxyFingerprint string
// MiddleEndFingerprint changes when restart-only [middle-end] settings
// change. It is a digest and never contains proxy credentials or tags.
MiddleEndFingerprint string
}
Config configures the gnet proxy server.
type ConnContext ¶
type ConnContext struct {
// contains filtered or unexported fields
}
ConnContext holds per-connection state for the gnet event handler.
func NewConnContext ¶
func NewConnContext() *ConnContext
NewConnContext creates a new connection context.
func (*ConnContext) Cleanup ¶ added in v0.1.8
func (c *ConnContext) Cleanup()
Cleanup zeros sensitive data in the connection context. Should be called when the connection is closed. Note: cipher.Stream internal state cannot be zeroed (opaque Go types).
func (*ConnContext) DCID ¶ added in v0.1.8
func (c *ConnContext) DCID() int
DCID returns the DC ID this connection is using (0 if not yet determined).
func (*ConnContext) ID ¶ added in v0.1.4
func (c *ConnContext) ID() uint64
ID returns the connection ID.
func (*ConnContext) LogPrefix ¶ added in v0.1.4
func (c *ConnContext) LogPrefix() string
LogPrefix returns a log prefix like "#123" or "#123:user1".
func (*ConnContext) ProtocolMode ¶ added in v0.3.3
func (c *ConnContext) ProtocolMode() ProtocolMode
ProtocolMode returns the protocol mode (ModeEE or ModeDD).
func (*ConnContext) RealClientAddr ¶ added in v0.1.4
func (c *ConnContext) RealClientAddr(fallback net.Addr) net.Addr
RealClientAddr returns the real client address from PROXY protocol. Falls back to the provided gnet connection's remote address if not set.
func (*ConnContext) Relay ¶
func (c *ConnContext) Relay() *RelayContext
Relay returns the relay context (lock-free, may be nil).
func (*ConnContext) SetProtocolMode ¶ added in v0.3.3
func (c *ConnContext) SetProtocolMode(mode ProtocolMode)
SetProtocolMode sets the protocol mode.
func (*ConnContext) SetRealClientAddr ¶ added in v0.1.4
func (c *ConnContext) SetRealClientAddr(addr net.Addr)
SetRealClientAddr sets the real client address from PROXY protocol.
func (*ConnContext) SetRelay ¶
func (c *ConnContext) SetRelay(r *RelayContext)
SetRelay sets the relay context and transitions to relay state.
func (*ConnContext) SetSpliceOverride ¶ added in v0.4.1
func (c *ConnContext) SetSpliceOverride(addr string)
SetSpliceOverride records a one-shot splice target ("host:port") for an unauthenticated probe whose SNI is on the mask safelist.
func (*ConnContext) SetState ¶
func (c *ConnContext) SetState(state ConnState)
SetState sets the connection state (lock-free).
func (*ConnContext) SetTrafficCounters ¶ added in v0.3.0
func (c *ConnContext) SetTrafficCounters(bytesIn, bytesOut *atomic.Int64)
SetTrafficCounters sets the traffic counter pointers for this connection. Must be called during handshake before entering relay state.
func (*ConnContext) State ¶
func (c *ConnContext) State() ConnState
State returns the current connection state (lock-free).
func (*ConnContext) TrafficIn ¶ added in v0.3.0
func (c *ConnContext) TrafficIn() *atomic.Int64
TrafficIn returns the traffic-in counter (may be nil). Safe to call without lock because SetTrafficCounters is called during handshake and getters are only called during relay state (after handshake).
func (*ConnContext) TrafficOut ¶ added in v0.3.0
func (c *ConnContext) TrafficOut() *atomic.Int64
TrafficOut returns the traffic-out counter (may be nil). Safe to call without lock because SetTrafficCounters is called during handshake and getters are only called during relay state (after handshake).
type ConnLimiter ¶ added in v0.1.4
type ConnLimiter struct {
// contains filtered or unexported fields
}
ConnLimiter limits concurrent connections per IP+secret combination. Uses sharded maps and atomic counters for minimal contention.
func NewConnLimiter ¶ added in v0.1.4
func NewConnLimiter(maxConns int) *ConnLimiter
NewConnLimiter creates a new connection limiter. maxConns <= 0 disables limiting.
func (*ConnLimiter) ActiveConnections ¶ added in v0.1.4
func (l *ConnLimiter) ActiveConnections() int64
ActiveConnections returns the total number of active connections being tracked. This is O(n) and should only be used for metrics, not in hot path.
func (*ConnLimiter) Release ¶ added in v0.1.4
func (l *ConnLimiter) Release(key string)
Release releases a connection slot. key must be the value returned by TryAcquire.
func (*ConnLimiter) TryAcquire ¶ added in v0.1.4
TryAcquire attempts to acquire a connection slot for the given IP+secret. Returns the key (for Release) and success status. If maxConns is 0, always succeeds (limiting disabled).
func (*ConnLimiter) TryAcquireIP ¶ added in v0.3.4
func (l *ConnLimiter) TryAcquireIP(ip net.IP) (key string, ok bool)
TryAcquireIP attempts to acquire a connection slot for the given IP only. Used in OnOpen to limit total connections per IP before authentication.
type ConnState ¶
type ConnState int32
ConnState represents the current state of a client connection.
const ( StateDetectProtocol ConnState = iota // Detect ee vs dd protocol from first bytes StateReadProxyProto // Need PROXY protocol header (optional) StateReadTLSHeader // Need 5 bytes for TLS record header StateReadTLSPayload // Need header.length bytes for payload StateReadO2Frame // Need 64 bytes for obfuscated2 frame (TLS-wrapped) StateReadDDFrame // Need 64 bytes for raw obfuscated2 frame (DD mode) StateDialingDC // Async dial in progress StateRelaying // Bidirectional relay active StateSplicing // Forward to mask host (invalid client) StateClosed // Connection is closing StateMiddleEnd // Authenticated client routed through Middle-End )
type DCConnContext ¶ added in v0.1.2
type DesyncDetector ¶ added in v0.1.8
type DesyncDetector struct {
// contains filtered or unexported fields
}
DesyncDetector tracks and reports protocol desynchronization events. Desync typically happens when crypto state diverges, causing decryption to produce garbage that looks like impossibly large frames.
func NewDesyncDetector ¶ added in v0.1.8
func NewDesyncDetector() *DesyncDetector
NewDesyncDetector creates a new desync detector.
func (*DesyncDetector) Report ¶ added in v0.1.8
func (d *DesyncDetector) Report( ctx *ConnContext, frameSize int, direction string, logger Logger, ) bool
Report reports a potential desync event. Returns true if this event was logged (not deduplicated). direction is "c2dc" (client to DC) or "dc2c" (DC to client).
type HandshakeFailureStat ¶ added in v0.5.2
HandshakeFailureStat is one cumulative MTProxy handshake failure counter.
type HotReloadConfig ¶ added in v0.1.8
type HotReloadConfig struct {
ConfigPath string // Path to config file
LoadConfig func() (*Config, string, error) // Config loader function
Handler *ProxyHandler
Logger Logger
SetLogFn func(level string) // Function to set log level
}
HotReloadConfig contains configuration for the hot reloader.
type HotReloader ¶ added in v0.1.8
type HotReloader struct {
// contains filtered or unexported fields
}
HotReloader watches a config file and reloads hot fields on change. Supports both file watching (fsnotify) and SIGHUP.
func NewHotReloader ¶ added in v0.1.8
func NewHotReloader(cfg HotReloadConfig) *HotReloader
NewHotReloader creates a new hot reloader.
func (*HotReloader) Start ¶ added in v0.1.8
func (r *HotReloader) Start()
Start begins watching for config changes. Returns immediately; watching runs in background goroutines.
func (*HotReloader) Stop ¶ added in v0.1.8
func (r *HotReloader) Stop()
Stop stops the hot reloader.
type InternalProxyAuth ¶ added in v0.5.0
type InternalProxyAuth struct {
// contains filtered or unexported fields
}
InternalProxyAuth authenticates the private in-process WEB-to-MTProxy hop. Its token is generated once, remains only in memory, and is never formatted.
func NewInternalProxyAuth ¶ added in v0.5.0
func NewInternalProxyAuth() (*InternalProxyAuth, error)
NewInternalProxyAuth creates one process-local authentication token.
func (*InternalProxyAuth) AppendPreface ¶ added in v0.5.0
func (a *InternalProxyAuth) AppendPreface(dst []byte) []byte
AppendPreface appends the private authentication preface to dst.
func (*InternalProxyAuth) GoString ¶ added in v0.5.0
func (*InternalProxyAuth) GoString() string
GoString prevents accidental token disclosure through Go-syntax formatting.
func (*InternalProxyAuth) String ¶ added in v0.5.0
func (*InternalProxyAuth) String() string
String prevents accidental token disclosure through ordinary formatting.
type Logger ¶
type Logger interface {
Debug(format string, args ...any)
Info(format string, args ...any)
Warn(format string, args ...any)
Error(format string, args ...any)
// DebugEnabled returns true if debug logging is enabled.
// Use to guard expensive debug log argument evaluation.
DebugEnabled() bool
}
Logger interface for proxy logging.
type LogicalQueueBudget ¶ added in v0.6.1
type LogicalQueueBudget struct {
Reserve func(bytes, items int) bool
Release func(bytes, items int)
}
LogicalQueueBudget charges retained payload capacity and queue items. The integration adds its standard item overhead. Callbacks must be thread-safe, must not call the stream, and must remain available until OnClosed returns. Existing ME decoder/request storage remains in the frontend's separate hard process budget; the logical transport does not duplicate that accounting.
type LogicalStream ¶ added in v0.6.1
type LogicalStream struct {
// contains filtered or unexported fields
}
LogicalStream is an MTProxy client driven by a real gnet owner loop. TryRead and TryWrite never block and run on that owner; (0, nil) means would block. Close may run anywhere. OnClosed acknowledges completed owner cleanup. For ME, this means detaching the client route and releasing logical queues; the shared manager still owns accepted requests and pending close controls.
func (*LogicalStream) Close ¶ added in v0.6.1
func (s *LogicalStream) Close() error
func (*LogicalStream) InboundBuffered ¶ added in v0.6.1
func (s *LogicalStream) InboundBuffered() int
func (*LogicalStream) LocalAddr ¶ added in v0.6.1
func (s *LogicalStream) LocalAddr() net.Addr
func (*LogicalStream) OutboundBuffered ¶ added in v0.6.1
func (s *LogicalStream) OutboundBuffered() int
func (*LogicalStream) OwnerStopped ¶ added in v0.6.1
func (s *LogicalStream) OwnerStopped()
OwnerStopped retires work accepted but abandoned by gnet during an unexpected exit. The integration must call it only after the owning engine has returned, when none of its callbacks can still run. Normal shutdown closes streams and waits for OnClosed before stopping the engine instead.
func (*LogicalStream) ReadableBytes ¶ added in v0.6.1
func (s *LogicalStream) ReadableBytes() int
ReadableBytes lets the owner reserve an exact destination before TryRead.
func (*LogicalStream) RemoteAddr ¶ added in v0.6.1
func (s *LogicalStream) RemoteAddr() net.Addr
func (*LogicalStream) TryRead ¶ added in v0.6.1
func (s *LogicalStream) TryRead(dst []byte) (int, error)
type LogicalStreamOptions ¶ added in v0.6.1
type LogicalStreamOptions struct {
Owner gnet.EventLoop
ClientAddr netip.AddrPort
LocalAddr net.Addr
// MaxInputBytes bounds parser staging (at most one item). Extra handshake
// and upstream relay storage is separately bounded by the core's record
// and relay limits, and charged to the same external InputBudget.
MaxInputBytes int
MaxOutputBytes int
MaxInputItems int
MaxOutputItems int
InputBudget LogicalQueueBudget
OutputBudget LogicalQueueBudget
Notify func()
OnOpened func(error)
OnClosed func(error)
}
LogicalStreamOptions supplies trusted process-local ingress metadata. The owner remains fixed across HTTP requests, carrier lanes, and reconnects.
type MiddleEndBindingSource ¶ added in v0.6.0
type MiddleEndBindingSource interface {
BindReady(middleend.DCID) (*middleend.ClientBinding, error)
Ready() <-chan struct{}
TryNextReady() *middleend.ClientReadyToken
Done() <-chan struct{}
}
MiddleEndBindingSource is the exclusive binding and readiness source for one frontend. A source may route through one fixed manager or supervise multiple whole generations. BindReady must select once, fix the ready-token consumer, and never migrate a returned binding. Each binding must expose the public source IP of its selected physical link. Ready and Done must return stable channels, and TryNextReady has one consumer.
type MiddleEndFrontendConfig ¶ added in v0.6.0
type MiddleEndFrontendConfig struct {
Source MiddleEndBindingSource
PrecommitFailure MiddleEndPrecommitAction
ProxyTag *middleend.ProxyTag
MaxPendingClientBytes int
MaxPendingClientBytesTotal int
MaxPendingOutputBytesTotal int
// ResponseBudget must be the same pool used by Source. Nil keeps legacy
// output accounting. Processing reserve must fit one complete encode.
ResponseBudget *middleend.ResponseBudget
OutputRetryInitial time.Duration
OutputRetryMax time.Duration
OutputStallTimeout time.Duration
}
MiddleEndFrontendConfig injects an already-started, externally owned binding source. Every operational field is required except ProxyTag. The frontend consumes Source readiness exclusively but never closes the source or its engine runtimes. Source must remain open until gnet has delivered every connection OnClose callback and the serving engine has returned.
func (MiddleEndFrontendConfig) GoString ¶ added in v0.6.0
func (c MiddleEndFrontendConfig) GoString() string
GoString redacts the binding source and optional process-wide proxy tag.
func (MiddleEndFrontendConfig) String ¶ added in v0.6.0
func (MiddleEndFrontendConfig) String() string
String redacts the binding source and optional process-wide proxy tag.
type MiddleEndFrontendStats ¶ added in v0.6.0
type MiddleEndFrontendStats struct {
MiddleEndBindingsActive int64
MiddleEndBindingsTotal uint64
DirectFallbacksActive int64
DirectFallbacksTotal uint64
InputBytes int64
InputBytesHighWater int64
InputBytesLimit int64
InputBackpressureEvents uint64
OutputBytes int64
OutputBytesHighWater int64
OutputBytesLimit int64
OutputBackpressureEvents uint64
OutputEvictions uint64
ResponseWaitsTotal [middleend.ResponseOutputWaitCount]uint64
ResponseWaitsCompletedTotal [middleend.ResponseOutputWaitCount]uint64
// Completed durations use microseconds to avoid nanosecond counter wrap
// after 21 days with 10,000 concurrent waits. Resolution is one microsecond.
ResponseWaitDurationMicroseconds [middleend.ResponseOutputWaitCount]uint64
ResponseWaitsActive [middleend.ResponseOutputWaitCount]int64
ResponseStallClosures uint64
}
MiddleEndFrontendStats reports frontend observations. OutputBytes counts unread bytes; the shared response pool separately charges retained capacity.
type MiddleEndPrecommitAction ¶ added in v0.6.0
type MiddleEndPrecommitAction uint8
MiddleEndPrecommitAction selects the only allowed outcome when Middle-End setup fails before BindReady succeeds. A successful binding is irreversible.
const ( MiddleEndPrecommitDirectFallback MiddleEndPrecommitAction = iota + 1 MiddleEndPrecommitClose )
type ProtocolMode ¶ added in v0.3.3
type ProtocolMode uint8
ProtocolMode indicates the MTProxy protocol variant.
const ( ModeEE ProtocolMode = iota // FakeTLS + Obfuscated2 (ee prefix) ModeDD // Raw Obfuscated2 (dd prefix) )
type ProxyHandler ¶
type ProxyHandler struct {
gnet.BuiltinEventEngine
// contains filtered or unexported fields
}
ProxyHandler implements gnet.EventHandler for the MTProxy server.
func NewProxyHandler ¶
func NewProxyHandler(cfg *Config, logger Logger) *ProxyHandler
NewProxyHandler creates a new gnet proxy handler.
func NewProxyHandlerWithMiddleEnd ¶ added in v0.6.0
func NewProxyHandlerWithMiddleEnd( cfg *Config, logger Logger, frontend MiddleEndFrontendConfig, ) (*ProxyHandler, error)
NewProxyHandlerWithMiddleEnd creates a handler with one explicit optional Middle-End branch. It constructs no link, engine, runtime, endpoint, or production configuration. Callers must start Source before accepting clients and must keep Source open until proxy shutdown and every connection OnClose callback have completed.
func RunWithHandler ¶ added in v0.1.8
func RunWithHandler(cfg *Config, logger Logger) (shutdown func(), handler *ProxyHandler, errCh <-chan error)
RunWithHandler starts the proxy and returns the handler for hot-reload. Returns a shutdown function, the handler (for hot-reload), and error channel.
func RunWithMiddleEnd ¶ added in v0.6.0
func RunWithMiddleEnd( cfg *Config, logger Logger, frontend MiddleEndFrontendConfig, ) (shutdown func(), handler *ProxyHandler, errCh <-chan error, err error)
RunWithMiddleEnd starts the proxy with an externally owned, already-started Middle-End binding source. It validates and constructs the handler before any network runtime starts. The caller must keep the source alive until shutdown returns and every public connection has received OnClose.
func (*ProxyHandler) ApplyHotConfig ¶ added in v0.1.8
func (h *ProxyHandler) ApplyHotConfig(cfg *Config)
ApplyHotConfig applies hot-reloadable configuration changes. Only certain fields can be changed at runtime; others require restart. Hot-reloadable fields:
- IdleTimeout: affects new connections only (existing keep their timeout)
Non-hot fields (require restart):
- BindAddr, Secrets, MaskHost/Port, ProxyProtocol, MaxIPsPerUser
func (*ProxyHandler) HandshakeFailureStats ¶ added in v0.5.2
func (h *ProxyHandler) HandshakeFailureStats() []HandshakeFailureStat
HandshakeFailureStats returns one stable counter for each processing stage.
func (*ProxyHandler) IdleTimeout ¶ added in v0.3.0
func (h *ProxyHandler) IdleTimeout() time.Duration
IdleTimeout returns the current idle timeout value (thread-safe).
func (*ProxyHandler) MiddleEndFrontendStats ¶ added in v0.6.0
func (h *ProxyHandler) MiddleEndFrontendStats() MiddleEndFrontendStats
MiddleEndFrontendStats returns aggregate public-side ME buffer pressure. A direct-only handler returns a zero snapshot.
func (*ProxyHandler) OnBoot ¶
func (h *ProxyHandler) OnBoot(eng gnet.Engine) gnet.Action
OnBoot is called when the gnet engine starts.
func (*ProxyHandler) OnShutdown ¶
func (h *ProxyHandler) OnShutdown(eng gnet.Engine)
OnShutdown is called when the gnet engine shuts down.
func (*ProxyHandler) OnTick ¶ added in v0.4.1
func (h *ProxyHandler) OnTick() (time.Duration, gnet.Action)
OnTick runs periodically (only when the client-silence breaker is enabled, via WithTicker) on event-loop 0. It closes relaying connections stuck in the iOS bad_salt wedge: the last relayed payload was server->client and the client has not answered for the configured threshold. A healthy connection whose last word was its own request/ack (lastClientByteMs >= lastServerByteMs) is never touched, and a connection where the client has not yet spoken in relay (lastClientByteMs == 0) is never touched. conn.Close() is cross-loop safe: it enqueues the close onto the connection's own event loop.
func (*ProxyHandler) OnTraffic ¶
func (h *ProxyHandler) OnTraffic(c gnet.Conn) gnet.Action
OnTraffic is called when data is available to read.
func (*ProxyHandler) OpenLogicalStream ¶ added in v0.6.1
func (h *ProxyHandler) OpenLogicalStream(options LogicalStreamOptions) (*LogicalStream, error)
func (*ProxyHandler) UserLimiter ¶ added in v0.3.0
func (h *ProxyHandler) UserLimiter() *UserIPLimiter
UserLimiter returns the user IP limiter for metrics access.
type ProxyProtoResult ¶ added in v0.1.4
type ProxyProtoResult struct {
SrcAddr net.Addr // Source (client) address
DstAddr net.Addr // Destination (server) address
HeaderLen int // Total bytes consumed by header
IsLocal bool // True if LOCAL command (health check)
}
ProxyProtoResult holds the parsed PROXY protocol header result.
func ParseProxyProtocol ¶ added in v0.1.4
func ParseProxyProtocol(data []byte) (*ProxyProtoResult, error)
ParseProxyProtocol attempts to parse a PROXY protocol v1 or v2 header. Returns nil if the data doesn't start with a valid PROXY protocol header. Returns error if header is malformed or incomplete.
type RelayContext ¶
type RelayContext struct {
// Client ciphers (client <-> proxy)
Encryptor cipher.Stream // encrypt data TO client
Decryptor cipher.Stream // decrypt data FROM client
// DC connection and ciphers (proxy <-> DC)
DCConn gnet.Conn // gnet connection to Telegram DC (enrolled in dcClient)
DCEncrypt cipher.Stream // encrypt data TO DC
DCDecrypt cipher.Stream // decrypt data FROM DC
ToDC *relayOutput // bounded cross-loop output to the DC owner
}
RelayContext holds immutable relay state set once after handshake. Read without locking via atomic pointer.
type ReplayCache ¶
type ReplayCache struct {
// contains filtered or unexported fields
}
ReplayCache detects replay attacks by tracking seen session IDs. Uses sharded TTL sets for:
- Proper LRU eviction (oldest entries removed first)
- TTL expiration checked during access, without background workers
- Reduced lock contention (64 shards)
func NewReplayCache ¶
func NewReplayCache(maxSize int, ttl time.Duration) *ReplayCache
NewReplayCache creates a new replay cache. maxSize is the total capacity across all shards. ttl is how long entries are kept. Nonpositive ttl disables expiration.
func (*ReplayCache) Len ¶ added in v0.1.8
func (c *ReplayCache) Len() int
Len returns the total number of entries across all shards.
func (*ReplayCache) Seen ¶
func (c *ReplayCache) Seen(sessionID []byte) bool
Seen checks if the session ID was seen before and adds it if not. Returns true if this is a replay attack, false if new. This operation is atomic (check-and-add).
type Secret ¶
type Secret struct {
Name string // User-friendly name for logging
Key []byte // 16-byte secret key
Host string // SNI hostname
RawHex string // Original hex string for link generation
}
Secret represents a named proxy secret.
type UserIPLimiter ¶ added in v0.3.0
type UserIPLimiter struct {
// contains filtered or unexported fields
}
UserIPLimiter tracks per-user statistics and optionally limits unique IPs per user. Uses sharded maps and LRU caches for minimal contention.
func NewUserIPLimiter ¶ added in v0.3.0
func NewUserIPLimiter(maxIPsPerUser int, blockTimeout time.Duration) *UserIPLimiter
NewUserIPLimiter creates a new user IP limiter/stats tracker. If maxIPsPerUser <= 0, limiting is disabled but stats are still tracked. Nonpositive blockTimeout keeps blocked IPs until capacity eviction or Close.
func (*UserIPLimiter) Close ¶ added in v0.3.0
func (l *UserIPLimiter) Close()
Close clears blocked IPs. The limiter owns no background workers.
func (*UserIPLimiter) LimitingEnabled ¶ added in v0.3.0
func (l *UserIPLimiter) LimitingEnabled() bool
LimitingEnabled returns whether IP limiting is active.
func (*UserIPLimiter) Release ¶ added in v0.3.0
func (l *UserIPLimiter) Release(key string)
Release releases a connection slot.
func (*UserIPLimiter) Stats ¶ added in v0.3.0
func (l *UserIPLimiter) Stats() []UserIPStats
Stats returns statistics for all users.
func (*UserIPLimiter) TrafficCounters ¶ added in v0.3.0
func (l *UserIPLimiter) TrafficCounters(secret []byte) (bytesIn, bytesOut *atomic.Int64)
TrafficCounters returns pointers to traffic counters for a user. Used to store in ConnContext for hot-path traffic counting.
func (*UserIPLimiter) TryAcquire ¶ added in v0.3.0
func (l *UserIPLimiter) TryAcquire(ip net.IP, secret []byte, secretName string) (key string, ok bool)
TryAcquire attempts to acquire a connection slot for the given IP+secret. Returns the key (for Release) and success status. secretName is used for metrics labeling. When limiting is disabled, always succeeds but still tracks stats.
type UserIPStats ¶ added in v0.3.0
type UserIPStats struct {
SecretName string
TrackedIPs int // Unique IPs in LRU cache (may include disconnected)
ActiveIPs int // IPs with active connections right now
BlockedIPs int
TrackedIPList []string // List of tracked IP addresses
BlockedIPList []string // List of currently blocked IP addresses
Connections int64
BytesIn int64
BytesOut int64
BlockedTotal int64
}
UserIPStats contains statistics for a single user.
Source Files
¶
- bufpool.go
- client_endpoint.go
- connlimit.go
- context.go
- dc_events.go
- dc_handler.go
- desync.go
- handler.go
- handshake.go
- hotreload.go
- idle_timeout.go
- internal_proxy_auth.go
- logical_output.go
- logical_stream.go
- middleend_budget.go
- middleend_frontend.go
- middleend_pressure.go
- middleend_response_output.go
- proxyproto.go
- relay.go
- relay_output.go
- replay.go
- runtime_stats.go
- server.go
- splice.go
- startup.go
- ttlset.go
- upstream_lifecycle.go
- userlimit.go
- web_diagnostics.go
- zeroize.go