httpHandler
# 1.入口
# 1.1.httpServer
tcp entrypoint初始化时会调用createHTTPServer()构造及启动http服务,http forwarder作为流量转发器注册到tcp router。// NewTCPEntryPoint creates a new TCPEntryPoint. func NewTCPEntryPoint(ctx context.Context, configuration *static.EntryPoint, hostResolverConfig *types.HostResolverConfig, openConnectionsGauge gokitmetrics.Gauge) (*TCPEntryPoint, error) { ... listener, err := buildListener(ctx, configuration) ... // 创建及启动http server httpServer, err := createHTTPServer(ctx, listener, configuration, true, reqDecorator) ... // http forwarder注册到tcp router rt.SetHTTPForwarder(httpServer.Forwarder) ... } func createHTTPServer(...) (*httpServer, error) { ... // 初始化http switcher httpSwitcher := middlewares.NewHandlerSwitcher(router.BuildDefaultHTTPRouter()) // 请求装饰器链包装switcher next, err := alice.New(requestdecorator.WrapHandler(reqDecorator)).Then(httpSwitcher) ... // 包装X-Forwarded-*检查器 // 代理服务器会接收前置代理的X-Forwarded-For或X-Forwarded-Proto形式的header // XForwarded handler会验证header是否可信(白名单的IP头),防止伪造IP handler, err = forwardedheaders.NewXForwarded( configuration.ForwardedHeaders.Insecure, configuration.ForwardedHeaders.TrustedIPs, next) ... // 附加修饰器(拒绝#...请求路径) handler = denyFragment(handler) // query参数;分隔符编码 if configuration.HTTP.EncodeQuerySemicolons { handler = encodeQuerySemicolons(handler) // 透传,接受附带;分隔符的查询字符串 } else { handler = http.AllowQuerySemicolons(handler) } // 禁止猜测Content-Type,避免误判导致安全问题 handler = contenttype.DisableAutoDetection(handler) // h2c if withH2c { // 包装,支持http/1.1和http/2 handler = h2c.NewHandler(handler, &http2.Server{ MaxConcurrentStreams: uint32(configuration.HTTP2.MaxConcurrentStreams), }) } ... // 配置keep-alive模式 if debugConnection || (configuration.Transport != nil && (configuration.Transport.KeepAliveMaxTime > 0 || configuration.Transport.KeepAliveMaxRequests > 0)) { // 包装keep-alive中间件,限制连接的存活时间和最大请求次数,防止长连接的资源泄漏 handler = newKeepAliveMiddleware(handler, configuration.Transport.KeepAliveMaxRequests, configuration.Transport.KeepAliveMaxTime) } // 构建http server serverHTTP := &http.Server{ Handler: handler, ErrorLog: stdlog.New(logs.NoLevel(log.Logger, zerolog.DebugLevel), "", 0), ReadTimeout: time.Duration(configuration.Transport.RespondingTimeouts.ReadTimeout), WriteTimeout: time.Duration(configuration.Transport.RespondingTimeouts.WriteTimeout), IdleTimeout: time.Duration(configuration.Transport.RespondingTimeouts.IdleTimeout), } ... if !strings.Contains(os.Getenv("GODEBUG"), "http2server=0") { // 启用http2 http2.ConfigureServer(serverHTTP, &http2.Server{ MaxConcurrentStreams: uint32(configuration.HTTP2.MaxConcurrentStreams), NewWriteScheduler: func() http2.WriteScheduler { return http2.NewPriorityWriteScheduler(nil) }, }) ... } // 包装tcp listener(用于tcp router推送流量) listener := (ln) go func() { // 启动http服务器 err := serverHTTP.Serve(listener) ... }() return &httpServer{ Server: serverHTTP, Forwarder: listener, Switcher: httpSwitcher, }, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96注意
createHTTPServer()会构造http switcher及挂载装饰器链,初始化http server及启动监听listener
# 1.2.http3Server
tcp entrypoint初始化会调用newHTTP3Server()构造http3服务,作为https服务的udp入口监听请求,复用https handler处理流量。func newHTTP3Server(...) (*http3server, error) { ... // 创建udp socket listenConfig := newListenConfig(configuration) conn, err := listenConfig.ListenPacket(ctx, "udp", configuration.GetAddress()) ... // 初始化http3服务器 h3 := &http3server{ http3conn: conn, getter: func(info *tls.ClientHelloInfo) (*tls.Config, error) { return nil, errors.New("no tls config") }, } // 构造server复用https handler h3.Server = &http3.Server{ Addr: configuration.GetAddress(), Port: configuration.HTTP3.AdvertisedPort, Handler: httpsServer.Server.(*http.Server).Handler, TLSConfig: &tls.Config{GetConfigForClient: h3.getGetConfigForClient}, } // 重写https的handler previousHandler := httpsServer.Server.(*http.Server).Handler // 其实就是流量处理时给个标识,告诉客户端支持h3 httpsServer.Server.(*http.Server).Handler = http.HandlerFunc(func(rw http.ResponseWriter,req *http.Request){ // 补充Alt-Svc: h3=":443";ma=2592000头信息 h3.Server.SetQuicHeaders(rw.Header()) ... // 基于原handler处理请求 previousHandler.ServeHTTP(rw, req) }) return h3, nil } // Start starts the TCP server. func (e *TCPEntryPoint) Start(ctx context.Context) { ... // 启动http3服务器 if e.http3Server != nil { go func() { _ = e.http3Server.Start() }() } ... }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49注意
http3 server监听udp listener,将请求转发到https handler,其实就是作为https服务的udp入口
# 2.路由
# 2.1.switchRouter
watcher监听到配置变更会推送到注册的listener,switchRouter会调用创建http handler及注册到http server。func switchRouter(...) func(conf dynamic.Configuration) { return func(conf dynamic.Configuration) { ... // 构造router routers, udpRouters := routerFactory.CreateRouters(rtConf) // 切换router及handler serverEntryPointsTCP.Switch(routers) ... } } // CreateRouters creates new TCPRouters and UDPRouters. func (f *RouterFactory) CreateRouters(rtConf *runtime.Configuration) (map[string]*tcprouter.Router, map[string]udp.Handler) { ... // HTTP routerManager := router.NewManager(rtConf, serviceManager, middlewaresBuilder, f.observabilityMgr, f.tlsManager) handlersNonTLS := routerManager.BuildHandlers(ctx, f.entryPointsTCP, false) handlersTLS := routerManager.BuildHandlers(ctx, f.entryPointsTCP, true) ... // TCP rtTCPManager := tcprouter.NewManager(rtConf, svcTCPManager, middlewaresTCPBuilder, handlersNonTLS, handlersTLS, f.tlsManager) // 构造的tcp router会挂载http handler routersTCP := rtTCPManager.BuildHandlers(ctx, f.entryPointsTCP) ... return routersTCP, routersUDP } // Switch the TCP routers. func (eps TCPEntryPoints) Switch(routersTCP map[string]*tcprouter.Router) { for entryPointName, rt := range routersTCP { eps[entryPointName].SwitchRouter(rt) } } // SwitchRouter switches the TCP router handler. func (e *TCPEntryPoint) SwitchRouter(rt *tcprouter.Router) { // 同步http forwarder rt.SetHTTPForwarder(e.httpServer.Forwarder) // 切换http handler httpHandler := rt.GetHTTPHandler() ... e.httpServer.Switcher.UpdateHandler(httpHandler) // 同步https forwarder rt.SetHTTPSForwarder(e.httpsServer.Forwarder) // 切换https handler httpsHandler := rt.GetHTTPSHandler() ... e.httpsServer.Switcher.UpdateHandler(httpsHandler) // 切换tcp router e.switcher.Switch(rt) // http3 handler切换 if e.http3Server != nil { e.http3Server.Switch(rt) } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67注意
http handler构造后会先挂到tcp router,后续利用switcher更新到http server
# 2.2.buildHandlers
buildHandlers根据entrypoint生成http muxer,http muxer挂到tcp router以供后续switcher到http server。// BuildHandlers Builds handler for all entry points. func (m *Manager) BuildHandlers(rootCtx context.Context,entryPoints []string,tls bool) map[string]http.Handler { ... // 定义路由入口 for entryPointName, routers := range m.getHTTPRouters(rootCtx, entryPoints, tls) { ... // 构造路由 handler, err := m.buildEntryPointHandler(ctx, entryPointName, routers) ... entryPointHandlers[entryPointName] = handler } // 无路由入口 for _, entryPointName := range entryPoints { ... handler, ok := entryPointHandlers[entryPointName] if ok || handler != nil { continue } // 构造默认路由 handler, err := m.observabilityMgr.BuildEPChain(ctx, entryPointName, "").Then(BuildDefaultHTTPRouter()) ... entryPointHandlers[entryPointName] = handler } return entryPointHandlers } // 构建router handler func (m *Manager) buildEntryPointHandler(ctx context.Context, entryPointName string, configs map[string]*runtime.RouterInfo) (http.Handler, error) { // 创建http muxer muxer, err := httpmuxer.NewMuxer() ... // 配置默认处理器(注入观测链,记录请求日志及指标) defaultHandler, err := m.observabilityMgr.BuildEPChain(ctx, entryPointName, "defaultHandler").Then(http.NotFoundHandler()) ... // NotFoundHandler加入muxer muxer.SetDefaultHandler(defaultHandler) // 遍历路由 for routerName, routerConfig := range configs { ... // 构建路由handler handler, err := m.buildRouterHandler(ctxRouter, routerName, routerConfig) ... // 挂上观测链(入口级) observabilityChain := m.observabilityMgr.BuildEPChain(ctx, entryPointName, routerConfig.Service) handler, err = observabilityChain.Then(handler) ... // 路由及规则加入muxer muxer.AddRoute(routerConfig.Rule, routerConfig.RuleSyntax, routerConfig.Priority, handler) ... } // 补充recover中间件链 chain := alice.New() chain = chain.Append(func(next http.Handler) (http.Handler, error) { return recovery.New(ctx, next) }) return chain.Then(muxer) } func (m *Manager) buildRouterHandler(ctx context.Context, routerName string, routerConfig *runtime.RouterInfo) (http.Handler, error) { // 缓存检查 if handler, ok := m.routerHandlers[routerName]; ok { return handler, nil } ... // 构建http handler handler, err := m.buildHTTPHandler(ctx, routerConfig, routerName) ... // router关联日志中间件(可选)及记录 handlerWithAccessLog, err := alice.New(func(next http.Handler) (http.Handler, error) { return accesslog.NewFieldHandler(next, accesslog.RouterName, routerName, nil), nil }).Then(handler) ... m.routerHandlers[routerName] = handlerWithAccessLog return m.routerHandlers[routerName], nil } func (m *Manager) buildHTTPHandler(ctx context.Context, router *runtime.RouterInfo, routerName string) (http.Handler, error) { ... // 收集middleware限定名 router.Middlewares = qualifiedNames ... // 构建service handler sHandler, err := m.serviceManager.BuildHTTP(ctx, router.Service) ... // 构建router级中间件链 mHandler := m.middlewaresBuilder.BuildChain(ctx, router.Middlewares) ... // 构建router级观测链 chain = chain.Append(observability.WrapMiddleware(ctx, metricsHandler)) ... // 默认router加一个中间件防止递归自己 if router.DefaultRule { chain = chain.Append(denyrouterrecursion.WrapHandler(routerName)) } ... // 组合中间件链和service handler return chain.Extend(*mHandler).Then(sHandler) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111注意
http muxer处理http流量本质也是基于路由匹配及执行http handler.ServeHTTP()
# 2.3.buildHTTP
serviceManager.BuildHTTP()会根据service创建handler,这个handler会挂到http muxer进一步注册到http server。// BuildHTTP Creates a http.Handler for a service configuration. func (m *Manager) BuildHTTP(rootCtx context.Context, serviceName string) (http.Handler, error) { serviceName = provider.GetQualifiedName(rootCtx, serviceName) ... handler, ok := m.services[serviceName] // 复用已存在的 if ok { return handler, nil } // 获取service配置 conf, ok := m.configs[serviceName] ... switch { // LB类型 case conf.LoadBalancer != nil: lb, err = m.getLoadBalancerServiceHandler(ctx, serviceName, conf) ... // 加权类型 case conf.Weighted != nil: lb, err = m.getWRRServiceHandler(ctx, serviceName, conf.Weighted) ... // 镜像类型 case conf.Mirroring != nil: lb, err = m.getMirrorServiceHandler(ctx, conf.Mirroring) ... // 故障转移类型 case conf.Failover != nil: lb, err = m.getFailoverServiceHandler(ctx, serviceName, conf.Failover) ... default: sErr := fmt.Errorf("the service %q does not have any type defined", serviceName) conf.AddError(sErr, true) return nil, sErr } m.services[serviceName] = lb return lb, nil } // LB类型处理(其他类型递归LB构造) func (m *Manager) getLoadBalancerServiceHandler(ctx context.Context, serviceName string, info *runtime.ServiceInfo) (http.Handler, error) { // 获取LB服务列表 service := info.LoadBalancer ... // 向连接池获取连接 roundTripper, err := m.roundTripperManager.Get(service.ServersTransport) ... // 初始化负载均衡器 lb := wrr.New(service.Sticky, service.HealthCheck != nil) ... // 打乱遍历后端 for _, server := range shuffle(service.Servers, m.rand) { // 基于FNV哈希生成唯一的proxyName hasher := fnv.New64a() _, _ = hasher.Write([]byte(server.URL)) // this will never return an error. proxyName := hex.EncodeToString(hasher.Sum(nil)) // 解析后端 target, err := url.Parse(server.URL) ... // 构建后端代理 proxy := buildSingleHostProxy(target, passHostHeader, time.Duration(flushInterval), roundTripper, m.bufferPool) ... // 后端代理加入负载均衡器 lb.Add(proxyName, proxy, server.Weight) // 更新后端状态可用 info.UpdateServerStatus(target.String(), runtime.StatusUp) healthCheckTargets[proxyName] = target } // LB开启健康检查 if service.HealthCheck != nil { // 构造hc检查器 m.healthCheckers[serviceName] = healthcheck.NewServiceHealthChecker(...) } return lb, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85注意
proxy可能指向externalName,也可能指向clusterIP,或直接指向endpoint
# 2.4.addRoute
http handler构造后会注册到http muxer,以结合matcher基于AST树匹配规则,满足规则的请求才转给service handler处理。// AddRoute add a new route to the router. func (m *Muxer) AddRoute(rule string, syntax string, priority int, handler http.Handler) error { ... switch syntax { case "v2": // 解析规则 parse, err = m.parserV2.Parse(rule) ... // 匹配回调 matcherFuncs = httpFuncsV2 default: parse, err = m.parser.Parse(rule) ... matcherFuncs = httpFuncs } // 生成匹配树 buildTree, ok := parse.(rules.TreeBuilder) ... // 规则及回调加入machiner err = matchers.addRule(buildTree(), matcherFuncs) ... // 注册muxer路由列表 m.routes = append(m.routes, &route{ handler: handler, matchers: matchers, priority: priority, }) // 根据优先级排序 sort.Sort(m.routes) return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
注意
router的核心就是matcher抽象规则树和service handler,前者匹配请求,后者反向代理真正服务,service handler本质上是LB
# 3.http流量
# 3.1.forwarder
tcp router分发http流量会执行forwarder.ServeHTTP()推送流量,http server会监听forwarder的推送管道以获取连接处理。// https专用(tcp router调用分发) func (t *TLSHandler) ServeTCP(conn WriteCloser) { t.Next.ServeTCP(tls.Server(conn, t.Config)) } // ServeTCP uses the connection to serve it later in "Accept". func (h *httpForwarder) ServeTCP(conn tcp.WriteCloser) { h.connChan <- conn } // Serve always returns a non-nil error and closes l. // After [Server.Shutdown] or [Server.Close], the returned error is [ErrServerClosed]. func (srv *Server) Serve(l net.Listener) error { ... for { // 推过来的conn rw, err := l.Accept() ... // 构造新连接 c := srv.newConn(rw) c.setState(c.rwc, StateNew, runHooks) // before Serve can return // 连接处理 go c.serve(connCtx) } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28注意
1.
httpForwarder负责将conn传给http server,http server基于conn新建http conn进行请求处理2.
https流量会基于tlsHandler转发请求到httpForwarder,与http流量处理的区别在于使用tls conn以解析流量
# 3.2.muxer
http muxer作为http handler挂到http server进行流量处理,muxter关联很多route,根据抽象语法树matcher匹配获取handler。// ServeHTTP forwards the connection to the matching HTTP handler. // Serves 404 if no handler is found. func (m *Muxer) ServeHTTP(rw http.ResponseWriter, req *http.Request) { // 匹配route handler for _, route := range m.routes { if route.matchers.match(req) { route.handler.ServeHTTP(rw, req) return } } // 默认404 handler m.defaultHandler.ServeHTTP(rw, req) }1
2
3
4
5
6
7
8
9
10
11
12
13
14注意
muxter匹配的handler类型包括balancer、mirror及failover,请求会转给handler.ServeHTTP()处理
# 3.3.balancer
balancer是http handler之一,接收到的请求流量最终会调用http handler处理,balancer是负载均衡器实现。// http.Handler的LB实现 func (b *Balancer) ServeHTTP(w http.ResponseWriter, req *http.Request) { // 客户端黏性 if b.stickyCookie != nil { // 代理名称 cookie, err := req.Cookie(b.stickyCookie.name) ... b.handlersMu.RLock() // 获取使用过的proxy handler, ok := b.handlerMap[cookie.Value] b.handlersMu.RUnlock() if ok && handler != nil { b.handlersMu.RLock() // 获取proxy状态 _, isHealthy := b.status[handler.name] b.handlersMu.RUnlock() // proxy健康 if isHealthy { // 请求处理 handler.ServeHTTP(w, req) return } } } // 获取proxy server, err := b.nextServer() ... // 客户端亲和 if b.stickyCookie != nil { // cookie记录proxy cookie := &http.Cookie{ Name: b.stickyCookie.name, Value: hash(server.name), Path: "/", HttpOnly: b.stickyCookie.httpOnly, Secure: b.stickyCookie.secure, SameSite: convertSameSite(b.stickyCookie.sameSite), MaxAge: b.stickyCookie.maxAge, } http.SetCookie(w, cookie) } // 请求处理 server.ServeHTTP(w, req) } func (b *Balancer) nextServer() (*namedHandler, error) { b.handlersMu.Lock() defer b.handlersMu.Unlock() ... var handler *namedHandler for { // 获取deadline最小的handler(堆顶) handler = heap.Pop(b).(*namedHandler) // 记录当前deadline b.curDeadline = handler.deadline // 基于权重更新deadline,权重越大越可能选中 handler.deadline += 1 / handler.weight // handler重入堆 heap.Push(b, handler) // 后端健康 if _, ok := b.status[handler.name]; ok { break } } return handler, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75补充
wrr本质上是balancer外部包了一层balancer实现service负载均衡,实现上是执行两次balancer.ServeHTTP()
# 3.4.mirror
mirror模式采用影子流量机制,内部还是balancer,真实的用户请求会复制到镜像服务,用户的响应仅来自主服务,镜像服务的请求响应会丢弃。// mirror handler func (m *Mirroring) ServeHTTP(rw http.ResponseWriter, req *http.Request) { // 获取镜像服务 mirrors := m.getActiveMirrors() if len(mirrors) == 0 { // 无镜像服务,直接转给主服务 m.handler.ServeHTTP(rw, req) return } // 可复用的请求体 rr, bytesRead, err := newReusableRequest(req, m.maxBodySize) ... // 主服务处理请求 m.handler.ServeHTTP(rw, rr.clone(req.Context())) ... // 异步处理镜像请求 m.routinePool.GoCtx(func(_ context.Context) { for _, handler := range mirrors { // 复制请求 r := rr.clone(req.Context()) // 分发流量给镜像服务 handler.ServeHTTP(m.rw, r.WithContext(contextStopPropagation{ctx})) } }) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29注意
mirror模式内的handler仍然是balancer,复用的是balancer.ServeHTTP()流量处理
# 3.5.failover
failover模式是一种故障转移机制,允许主服务异常时自动切换到下一个备用服务,以维持请求的可用性,规避客户端的重试。// failover handler func (f *Failover) ServeHTTP(w http.ResponseWriter, req *http.Request) { f.handlerStatusMu.RLock() // 获取主服务状态 handlerStatus := f.handlerStatus f.handlerStatusMu.RUnlock() // 主服务健康 if handlerStatus { // 流量转向主服务 f.handler.ServeHTTP(w, req) return } f.fallbackStatusMu.RLock() // 获取备服务状态 fallbackStatus := f.fallbackStatus f.fallbackStatusMu.RUnlock() // 备服务健康 if fallbackStatus { // 流量转向备服务 f.fallbackHandler.ServeHTTP(w, req) return } // 响应服务异常 http.Error(w, http.StatusText(http.StatusServiceUnavailable), http.StatusServiceUnavailable) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29注意
failover的handler和fallback handler本质上也是balancer