Merge remote-tracking branch 'upstream/main' into merge-gramsrv-9106877

This commit is contained in:
onysd 2026-08-03 23:29:20 +03:00
commit ac6a50c5ff
697 changed files with 100880 additions and 8052 deletions

View file

@ -56,7 +56,7 @@ type legacyRPCHandlerWithMethod interface {
// LayerRPCHandler is the production API-RPC boundary. Admission is a separate
// allocation-bounded phase so the edge can freeze the connection profile,
// validate wrapper dependencies and establish exact request identity before
// flight/cache/scheduler ownership is acquired.
// execution-ledger/scheduler ownership is acquired.
type LayerRPCHandler interface {
AdmitLayer(profile tlprofile.Profile, b *bin.Buffer, limits tlprofile.Limits) (tlprofile.Admission, error)
AdmitUnprofiled(b *bin.Buffer, limits tlprofile.Limits) (tlprofile.Admission, error)
@ -70,6 +70,14 @@ type LayerRPCHandler interface {
) (tlprofile.Result, string, error)
}
// LayerRPCOptionsAdmitter extends the stable handler boundary with
// caller-owned admission capabilities. Implementations that do not expose it
// remain usable for requests that need only Limits.
type LayerRPCOptionsAdmitter interface {
AdmitLayerWithOptions(profile tlprofile.Profile, b *bin.Buffer, options tlprofile.AdmissionOptions) (tlprofile.Admission, error)
AdmitUnprofiledWithOptions(b *bin.Buffer, options tlprofile.AdmissionOptions) (tlprofile.Admission, error)
}
// LayerRPCDefaultProfileAdmitter decodes with a recoverable inherited/default
// profile. Production handlers should implement it with the same sparse
// tlprofile dispatcher and semantic adapter registry used by AdmitLayer. The split keeps old
@ -79,6 +87,23 @@ type LayerRPCDefaultProfileAdmitter interface {
AdmitDefaultLayer(profile tlprofile.Profile, b *bin.Buffer, limits tlprofile.Limits) (tlprofile.Admission, error)
}
// LayerRPCDefaultProfileOptionsAdmitter is the capability-aware form used when
// exact admission needs caller-owned resources such as bounded gzip expansion.
type LayerRPCDefaultProfileOptionsAdmitter interface {
AdmitDefaultLayerWithOptions(profile tlprofile.Profile, b *bin.Buffer, options tlprofile.AdmissionOptions) (tlprofile.Admission, error)
}
// LayerRPCFlatBytesPayloadSizer is an optional, allocation-free admission
// capability for exact terminal requests whose generated object graph contains
// one already-bounded flat bytes payload. The handler may return ok only after
// proving the complete terminal wire shape and every semantic field cap it
// relies on. The edge still owns the multiplier, graph slack and all process /
// connection budgets; an absent or invalid hint falls back to the conservative
// generic graph charge.
type LayerRPCFlatBytesPayloadSizer interface {
LayerRPCFlatBytesPayloadSize(wire []byte) (payloadBytes int, ok bool)
}
// LayerRPCSessionProfileResolver may restore an exact profile only when it was
// previously proven for this same (auth_key_id, session_id). Auth-key-wide
// device metadata is intentionally ineligible: a client upgrade can reuse its
@ -89,7 +114,7 @@ type LayerRPCSessionProfileResolver interface {
// LayerRPCOrderedSessionProfileResolver restores both the selected Layer and
// the newest invokeWithLayer client msg_id which proved it. The cursor prevents
// an old cached request replay on a replacement physical connection from
// an old retained request replay on a replacement physical connection from
// rolling the logical session back to an older profile.
type LayerRPCOrderedSessionProfileResolver interface {
NegotiatedSessionLayerEvidence(authKeyID [8]byte, sessionID int64) (layer int, msgID int64, ok bool)
@ -161,7 +186,7 @@ type LayerRPCDurableSessionProfileDeleter interface {
// LayerRPCReplayPreparer reapplies connection-local wrapper state for an
// already-executed exact request without consuming its one-shot business
// dispatch lease. The returned callback is safe to run only after a successful
// cached rpc_result reaches the replacement physical connection.
// replayed rpc_result reaches the replacement physical connection.
type LayerRPCReplayPreparer interface {
PrepareAdmittedReplay(
ctx context.Context,
@ -187,7 +212,7 @@ type LayerRPCProfileEvidenceContext interface {
// LayerRPCAdmissionProfilePublisher advances the auth-key-wide inherited
// default for fresh explicit evidence. admissionSeq is allocated once by the
// edge's exact flight owner and globally orders different MTProto sessions;
// cached joins/replays never call this hook again.
// joined/replayed requests never call this hook again.
type LayerRPCAdmissionProfilePublisher interface {
PublishAdmittedLayerProfileEvidence(
ctx context.Context,
@ -221,12 +246,14 @@ type Options struct {
// codec 必须是 gotd 内置四种 codec可包 NoHeader或实现 InboundFrameBudgetedCodec
// 无法在 payload 分配前预检长度的 codec 会 fail-closed。
Codec func() transport.Codec
// ObfuscatedTCP 先按 MTProto TCP obfuscation 解包,再自动探测 codec。
// Telegram Desktop 的 tcpo_only endpoint 会走这个 64 字节前缀流程。
// ObfuscatedTCP 允许裸 TCP 使用 MTProto transport obfuscation。开启时按每条
// 物理连接的首 1/4/8 字节自动区分明文 transport 与 64-byte obfuscated2
// 随后把 wire mode + codec 冻结到同一条双向连接Telegram Desktop 的
// tcpo_only 与未开启混淆的第三方客户端可共用同一端口。
ObfuscatedTCP bool
// WebSocket 在同一个 listener 上接受 MTProto over WebSocket(/apiws*)。
// 开启后仅在连接建立时读取前 4 字节做 HTTP/TCP 分流MTProto TCP
// 后续仍走原 ObfuscatedTCP + codec 热路径。
// 后续仍走原 TCP wire-mode + codec 探测路径。
WebSocket bool
// WebSocketAllowedOrigins 是允许浏览器发起 WebSocket upgrade 的页面 origin。
// 空列表表示只接受无 Origin 的非浏览器客户端;"*" 表示允许所有来源(仅调试)。
@ -268,19 +295,16 @@ type Options struct {
// 等于 copied bodyexact charge 是 typed decode 前的保守 materialization
// 上界,因此该配置不表示可并发接收 512 MiB wire body。默认 512 MiB。
RPCGlobalMaxBytes int64
// RPCResultCache* limits bound pending ownership and completed rpc_result
// replay state across the full 331-second duplicate horizon. Every owner is
// charged simultaneously at global, raw-auth and session scopes. Defaults:
// global 262144/64 MiB, auth 32768/32 MiB, session 16384/16 MiB.
RPCResultCacheMaxEntries int
RPCResultCacheMaxBytes int64
RPCResultCacheAuthMaxEntries int
RPCResultCacheAuthMaxBytes int64
RPCResultCacheSessionMaxEntries int
RPCResultCacheSessionMaxBytes int64
// RPCResultPendingPerAuth is an additional active-owner bound, independent
// RPCExecution*Entries bound in-flight owners and compact completed
// receipts. Payload bytes are not charged here: the logical-session
// outbox owns them under OutboundTrackedGlobalMaxBytes until ACK. ACK removes
// the receipt immediately; 331 seconds is only the no-ACK safety horizon.
RPCExecutionMaxEntries int
RPCExecutionAuthMaxEntries int
RPCExecutionSessionMaxEntries int
// RPCExecutionPendingPerAuth is an additional active-owner bound, independent
// from the retained entry limits and RPCGlobalMaxTasks. Default 2048.
RPCResultPendingPerAuth int
RPCExecutionPendingPerAuth int
// InboundFrameGlobalMaxBytes 是所有物理连接当前正在处理的 transport wire buffer
// 与最大解密 plaintext buffer 的总预算。长度前缀读取后、payload 分配前预留,默认
// 512 MiB非正值使用默认值。
@ -328,7 +352,7 @@ type Options struct {
// generated Layer admission by configuring the canonical-only route.
legacyRPC legacyRPCHandler
// LayerRPC is the generated exact-profile production path. When configured,
// every API request must complete admission before flight/cache scheduling.
// every API request must complete admission before execution-ledger scheduling.
LayerRPC LayerRPCHandler
// Metrics 接收连接层指标。默认 NopMetrics。
Metrics Metrics
@ -385,28 +409,19 @@ func (o *Options) setDefaults() {
if o.RPCGlobalMaxBytes <= 0 {
o.RPCGlobalMaxBytes = 512 << 20
}
if o.RPCResultCacheMaxEntries == 0 {
o.RPCResultCacheMaxEntries = rpcResultCacheMaxEntries
if o.RPCExecutionMaxEntries == 0 {
o.RPCExecutionMaxEntries = rpcExecutionMaxEntries
}
if o.RPCResultCacheMaxBytes == 0 {
o.RPCResultCacheMaxBytes = rpcResultCacheMaxBytes
if o.RPCExecutionAuthMaxEntries == 0 {
o.RPCExecutionAuthMaxEntries = rpcExecutionAuthMaxEntries
}
if o.RPCResultCacheAuthMaxEntries == 0 {
o.RPCResultCacheAuthMaxEntries = rpcResultCacheAuthMaxEntries
if o.RPCExecutionSessionMaxEntries == 0 {
o.RPCExecutionSessionMaxEntries = rpcExecutionSessionMaxEntries
}
if o.RPCResultCacheAuthMaxBytes == 0 {
o.RPCResultCacheAuthMaxBytes = rpcResultCacheAuthMaxBytes
}
if o.RPCResultCacheSessionMaxEntries == 0 {
o.RPCResultCacheSessionMaxEntries = rpcResultCacheSessionMaxEntries
}
if o.RPCResultCacheSessionMaxBytes == 0 {
o.RPCResultCacheSessionMaxBytes = rpcResultCacheSessionMaxBytes
}
if o.RPCResultPendingPerAuth == 0 {
o.RPCResultPendingPerAuth = rpcResultFlightMaxPendingPerAuth
if o.RPCResultPendingPerAuth > o.RPCGlobalMaxTasks {
o.RPCResultPendingPerAuth = o.RPCGlobalMaxTasks
if o.RPCExecutionPendingPerAuth == 0 {
o.RPCExecutionPendingPerAuth = rpcExecutionPendingPerAuth
if o.RPCExecutionPendingPerAuth > o.RPCGlobalMaxTasks {
o.RPCExecutionPendingPerAuth = o.RPCGlobalMaxTasks
}
}
if o.InboundFrameGlobalMaxBytes <= 0 {
@ -441,30 +456,19 @@ func (o *Options) setDefaults() {
}
}
func validateRPCResultCacheOptions(o Options) error {
if o.RPCResultCacheMaxEntries <= 0 || o.RPCResultCacheAuthMaxEntries <= 0 || o.RPCResultCacheSessionMaxEntries <= 0 {
return fmt.Errorf("rpc_result cache entry limits must be positive")
func validateRPCExecutionOptions(o Options) error {
if o.RPCExecutionMaxEntries <= 0 || o.RPCExecutionAuthMaxEntries <= 0 || o.RPCExecutionSessionMaxEntries <= 0 {
return fmt.Errorf("rpc execution ledger entry limits must be positive")
}
if o.RPCResultCacheMaxEntries < o.RPCResultCacheAuthMaxEntries ||
o.RPCResultCacheAuthMaxEntries < o.RPCResultCacheSessionMaxEntries {
return fmt.Errorf("rpc_result cache entry hierarchy must satisfy global >= auth >= session: %d/%d/%d",
o.RPCResultCacheMaxEntries, o.RPCResultCacheAuthMaxEntries, o.RPCResultCacheSessionMaxEntries)
if o.RPCExecutionMaxEntries < o.RPCExecutionAuthMaxEntries ||
o.RPCExecutionAuthMaxEntries < o.RPCExecutionSessionMaxEntries {
return fmt.Errorf("rpc execution ledger entry hierarchy must satisfy global >= auth >= session: %d/%d/%d",
o.RPCExecutionMaxEntries, o.RPCExecutionAuthMaxEntries, o.RPCExecutionSessionMaxEntries)
}
if o.RPCResultCacheMaxBytes < int64(maxOutboundBodyBytes) ||
o.RPCResultCacheAuthMaxBytes < int64(maxOutboundBodyBytes) ||
o.RPCResultCacheSessionMaxBytes < int64(maxOutboundBodyBytes) {
return fmt.Errorf("rpc_result cache byte limits must each be at least max outbound body %d: %d/%d/%d",
maxOutboundBodyBytes, o.RPCResultCacheMaxBytes, o.RPCResultCacheAuthMaxBytes, o.RPCResultCacheSessionMaxBytes)
}
if o.RPCResultCacheMaxBytes < o.RPCResultCacheAuthMaxBytes ||
o.RPCResultCacheAuthMaxBytes < o.RPCResultCacheSessionMaxBytes {
return fmt.Errorf("rpc_result cache byte hierarchy must satisfy global >= auth >= session: %d/%d/%d",
o.RPCResultCacheMaxBytes, o.RPCResultCacheAuthMaxBytes, o.RPCResultCacheSessionMaxBytes)
}
if o.RPCResultPendingPerAuth <= 0 || o.RPCResultPendingPerAuth > o.RPCGlobalMaxTasks ||
o.RPCResultPendingPerAuth > o.RPCResultCacheAuthMaxEntries {
return fmt.Errorf("rpc_result per-auth pending limit %d must be positive and <= global pending %d and auth entries %d",
o.RPCResultPendingPerAuth, o.RPCGlobalMaxTasks, o.RPCResultCacheAuthMaxEntries)
if o.RPCExecutionPendingPerAuth <= 0 || o.RPCExecutionPendingPerAuth > o.RPCGlobalMaxTasks ||
o.RPCExecutionPendingPerAuth > o.RPCExecutionAuthMaxEntries {
return fmt.Errorf("rpc execution per-auth pending limit %d must be positive and <= global pending %d and auth entries %d",
o.RPCExecutionPendingPerAuth, o.RPCGlobalMaxTasks, o.RPCExecutionAuthMaxEntries)
}
return nil
}
@ -510,7 +514,7 @@ type Server struct {
types *tmap.Map
admission *admissionController
rpcResults *rpcResultCache
rpcResults *rpcExecutionLedger
rpcRewrap *rpcRewrapRegistry
// onFrame 是测试钩子:收到一帧时回调其字节数;生产为 nil。
@ -520,14 +524,14 @@ type Server struct {
// New 创建 Server。
func New(opts Options) *Server {
opts.setDefaults()
if err := validateRPCResultCacheOptions(opts); err != nil {
panic(fmt.Sprintf("mtprotoedge: invalid result-cache options: %v", err))
if err := validateRPCExecutionOptions(opts); err != nil {
panic(fmt.Sprintf("mtprotoedge: invalid rpc execution options: %v", err))
}
conns := opts.ActiveSessions
if conns == nil {
conns = NewSessionManager(opts.Logger.Named("sessions"))
}
return &Server{
server := &Server{
log: opts.Logger,
codec: opts.Codec,
obfuscated: opts.ObfuscatedTCP,
@ -560,19 +564,21 @@ func New(opts Options) *Server {
clock: opts.Clock,
rand: opts.Rand,
types: tmap.New(tg.TypesMap(), mt.TypesMap(), proto.TypesMap()),
rpcResults: newRPCResultCacheWithFairCapacity(opts.Clock.Now, rpcResultCacheCapacity{
rpcResults: newRPCExecutionLedger(opts.Clock.Now, rpcExecutionLedgerCapacity{
maxPending: opts.RPCGlobalMaxTasks,
maxPendingPerAuth: opts.RPCResultPendingPerAuth,
globalMaxBytes: opts.RPCResultCacheMaxBytes,
globalMaxEntries: opts.RPCResultCacheMaxEntries,
authMaxBytes: opts.RPCResultCacheAuthMaxBytes,
authMaxEntries: opts.RPCResultCacheAuthMaxEntries,
sessionMaxBytes: opts.RPCResultCacheSessionMaxBytes,
sessionMaxEntries: opts.RPCResultCacheSessionMaxEntries,
maxPendingPerAuth: opts.RPCExecutionPendingPerAuth,
globalMaxEntries: opts.RPCExecutionMaxEntries,
authMaxEntries: opts.RPCExecutionAuthMaxEntries,
sessionMaxEntries: opts.RPCExecutionSessionMaxEntries,
replayStore: conns,
}),
rpcRewrap: newRPCRewrapRegistry(opts.RPCGlobalMaxTasks),
admission: newAdmissionController(opts.MaxConnections, opts.MaxConnectionsPerIP, opts.MaxConcurrentHandshakes),
}
conns.setLogicalSessionReleaseHook(func(key sessionKey) {
server.rpcResults.forgetSession(key.authKeyID, key.sessionID)
})
return server
}
// ListenAndServe binds the public MTProto socket and immediately enters Serve.
@ -634,8 +640,16 @@ func (s *Server) buildConn(tc transport.Conn, lease *physicalTransportLease, key
outboundTrackedBudget: s.outboundTrackedBudget,
outboundControlTrackedBudget: s.outboundControlBudget,
outboundScratchPool: s.outboundScratchPool,
rpcResultAcked: s.rpcRewrap.acknowledge,
rpcResultAcked: func(conn *Conn, reqMsgID int64) {
// The sole outbound actor invokes this only after resolving a client
// msgs_ack server msg_id through its tracked resend frame. The actor has
// already removed the sole outbox frame; now delete its receipt and
// retire any init-rewrap bookkeeping.
s.rpcResults.Acknowledge(conn.authKeyID, conn.sessionID, reqMsgID)
s.rpcRewrap.acknowledge(conn, reqMsgID)
},
}
s.conns.attachLogicalSession(c, s.outboundTrackedBudget)
c.startOutbound()
c.startInboundRPCScheduler(s.rpcScheduler, s.rpcInflight, s.rpcQueueSize, s.rpcTimeout)
return c
@ -648,6 +662,7 @@ func (s *Server) Serve(ctx context.Context, ln net.Listener) error {
// serveTCP/serveMixed 返回前会等待连接 goroutine 收敛,各 Conn 已先排空/取消任务;
// 最后再停止全局池,避免关闭过程中留下无人消费但仍占预算的队列。
s.rpcScheduler.start()
defer s.conns.releaseAllLogicalSessions()
defer s.rpcScheduler.stop(rpcCloseWaitTimeout)
// 只在最外层 listener 包一次,确保 same-port mux 的 sniff/HTTP upgrade 也计入
// raw admission而不是等连接已经分流后才计数。
@ -662,7 +677,11 @@ func (s *Server) serveTCP(ctx context.Context, ln net.Listener) error {
ctx, cancel := context.WithCancel(ctx)
defer cancel()
s.log.Info("Serving", zap.String("addr", ln.Addr().String()), zap.Int("dc", s.dc), zap.Bool("obfuscated_tcp", s.obfuscated))
s.log.Info("Serving",
zap.String("addr", ln.Addr().String()),
zap.Int("dc", s.dc),
zap.String("tcp_transport_mode", intakeTransport(s.obfuscated)),
)
defer s.log.Info("Stopped")
errCh := make(chan error, 1)
go func() {
@ -698,7 +717,7 @@ func (s *Server) serveMixed(ctx context.Context, ln net.Listener) error {
s.log.Info("Serving",
zap.String("addr", ln.Addr().String()),
zap.Int("dc", s.dc),
zap.Bool("obfuscated_tcp", s.obfuscated),
zap.String("tcp_transport_mode", intakeTransport(s.obfuscated)),
zap.Bool("websocket", true),
zap.Strings("websocket_origins", s.websocketOrigins),
)
@ -723,7 +742,7 @@ func (s *Server) serveMixed(ctx context.Context, ln net.Listener) error {
defer wg.Done()
errCh <- mux.Serve(ctx)
}()
// 裸 MTProto TCP每条连接在自己的 goroutine 里完成去混淆 + codec 探测。
// 裸 MTProto TCP每条连接在自己的 goroutine 里完成 wire mode + codec 探测。
go func() {
defer wg.Done()
errCh <- s.acceptLoop(ctx, mux.TCP(), s.obfuscated)
@ -764,11 +783,11 @@ func (s *Server) serveMixed(ctx context.Context, ln net.Listener) error {
return firstErr
}
// acceptLoop 接受裸连接,并为每条连接单独起 goroutine 完成「去混淆 + codec 探测 +
// acceptLoop 接受裸连接,并为每条连接单独起 goroutine 完成「wire mode + codec 探测 +
// serveConn」。探测在 accept 循环之外、带握手超时进行——慢/半开/坏 init 的客户端只占用
// 自己的 goroutine绝不阻塞其他连接的接入单条连接的握手失败也只关闭该连接不会拖垮
// 整个监听循环。obfuscated 为 true 时先走 obfuscated2 去混淆(裸 MTProto TCPWebSocket
// 连接传 falsegotd 升级处理器已完成去混淆)。
// 整个监听循环。obfuscated 为 true 时自动区分 plain 与 obfuscated2裸 MTProto TCP
// WebSocket 连接传 falsegotd 升级处理器已完成去混淆)。
func (s *Server) acceptLoop(ctx context.Context, ln net.Listener, obfuscated bool) error {
return s.acceptLoopTransport(ctx, ln, obfuscated, intakeTransport(obfuscated))
}
@ -829,7 +848,7 @@ func (s *Server) acceptLoopTransport(ctx context.Context, ln net.Listener, obfus
func (s *Server) serveDetectedConn(ctx context.Context, raw net.Conn, obfuscated bool, transportName string) {
started := time.Now()
remote, local := connRemote(raw), connLocal(raw)
// 握手读超时只覆盖去混淆 + codec 探测这一小段用真实墙钟时间SetReadDeadline 语义),
// 握手读超时只覆盖 wire-mode + codec 探测这一小段用真实墙钟时间SetReadDeadline 语义),
// 不走可能被测试注入的逻辑 clock。
if err := raw.SetReadDeadline(time.Now().Add(s.handshakeTimeout)); err != nil {
_ = raw.Close()
@ -848,7 +867,10 @@ func (s *Server) serveDetectedConn(ctx context.Context, raw net.Conn, obfuscated
}
}()
conn, err := s.promoteConn(raw, obfuscated)
conn, detectedTransport, err := s.promoteConn(raw, obfuscated)
if detectedTransport != "" {
transportName = detectedTransport
}
close(promoted)
if err != nil {
outcome := "error"
@ -885,15 +907,15 @@ func (s *Server) serveDetectedConn(ctx context.Context, raw net.Conn, obfuscated
}
}
// promoteConn 复用与 listener 组合完全一致的「obfuscated2 去混淆 + codec 探测」管线,但针对
// 单条连接,使其可在 accept 循环之外执行。obfuscated 对 WebSocket 连接必须为 falsegotd
// 升级处理器已剥离 obfuscated2 并补回 codec tag
func (s *Server) promoteConn(raw net.Conn, obfuscated bool) (transport.Conn, error) {
var ln net.Listener = newSingleConnListener(raw)
// promoteConn 针对单条连接执行一次 transport 提升。obfuscated=true 表示裸 TCP
// 允许混淆并自动区分 plain/obfuscated2WebSocket 必须传 false因为 gotd upgrade
// handler 已剥离 obfuscated2 并补回 codec tag。
func (s *Server) promoteConn(raw net.Conn, obfuscated bool) (transport.Conn, string, error) {
if obfuscated {
ln = transport.ObfuscatedListener(ln)
return s.promoteMixedTCP(raw)
}
return newCompatTransportListener(s.codec, ln, s.frameBudget).Accept()
conn, err := newCompatTransportConn(s.codec, raw, s.frameBudget)
return conn, "", err
}
// serveConn 处理单个传输连接:读帧并按 auth_key_id 分流。