()
| 1096 | } |
| 1097 | |
| 1098 | func (n *Node) startRPC() ([]net.Listener, error) { |
| 1099 | err := n.ConfigureRPC() |
| 1100 | if err != nil { |
| 1101 | return nil, err |
| 1102 | } |
| 1103 | |
| 1104 | listenAddrs := splitAndTrimEmpty(n.config.RPC.ListenAddress, ",", " ") |
| 1105 | |
| 1106 | if n.config.RPC.Unsafe { |
| 1107 | rpccore.AddUnsafeRoutes() |
| 1108 | } |
| 1109 | |
| 1110 | config := rpcserver.DefaultConfig() |
| 1111 | config.MaxBodyBytes = n.config.RPC.MaxBodyBytes |
| 1112 | config.MaxHeaderBytes = n.config.RPC.MaxHeaderBytes |
| 1113 | config.MaxOpenConnections = n.config.RPC.MaxOpenConnections |
| 1114 | // If necessary adjust global WriteTimeout to ensure it's greater than |
| 1115 | // TimeoutBroadcastTxCommit. |
| 1116 | // See https://github.com/DeAI-Artist/Linkis/issues/3435 |
| 1117 | if config.WriteTimeout <= n.config.RPC.TimeoutBroadcastTxCommit { |
| 1118 | config.WriteTimeout = n.config.RPC.TimeoutBroadcastTxCommit + 1*time.Second |
| 1119 | } |
| 1120 | |
| 1121 | // we may expose the rpc over both a unix and tcp socket |
| 1122 | listeners := make([]net.Listener, len(listenAddrs)) |
| 1123 | for i, listenAddr := range listenAddrs { |
| 1124 | mux := http.NewServeMux() |
| 1125 | rpcLogger := n.Logger.With("module", "rpc-server") |
| 1126 | wmLogger := rpcLogger.With("protocol", "websocket") |
| 1127 | wm := rpcserver.NewWebsocketManager(rpccore.Routes, |
| 1128 | rpcserver.OnDisconnect(func(remoteAddr string) { |
| 1129 | err := n.eventBus.UnsubscribeAll(context.Background(), remoteAddr) |
| 1130 | if err != nil && err != tmpubsub.ErrSubscriptionNotFound { |
| 1131 | wmLogger.Error("Failed to unsubscribe addr from events", "addr", remoteAddr, "err", err) |
| 1132 | } |
| 1133 | }), |
| 1134 | rpcserver.ReadLimit(config.MaxBodyBytes), |
| 1135 | ) |
| 1136 | wm.SetLogger(wmLogger) |
| 1137 | mux.HandleFunc("/websocket", wm.WebsocketHandler) |
| 1138 | rpcserver.RegisterRPCFuncs(mux, rpccore.Routes, rpcLogger) |
| 1139 | listener, err := rpcserver.Listen( |
| 1140 | listenAddr, |
| 1141 | config, |
| 1142 | ) |
| 1143 | if err != nil { |
| 1144 | return nil, err |
| 1145 | } |
| 1146 | |
| 1147 | var rootHandler http.Handler = mux |
| 1148 | if n.config.RPC.IsCorsEnabled() { |
| 1149 | corsMiddleware := cors.New(cors.Options{ |
| 1150 | AllowedOrigins: n.config.RPC.CORSAllowedOrigins, |
| 1151 | AllowedMethods: n.config.RPC.CORSAllowedMethods, |
| 1152 | AllowedHeaders: n.config.RPC.CORSAllowedHeaders, |
| 1153 | }) |
| 1154 | rootHandler = corsMiddleware.Handler(mux) |
| 1155 | } |
no test coverage detected