api_request.go 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204
  1. package channel
  2. import (
  3. "context"
  4. "errors"
  5. "fmt"
  6. "io"
  7. "net/http"
  8. common2 "one-api/common"
  9. "one-api/relay/common"
  10. "one-api/relay/constant"
  11. "one-api/relay/helper"
  12. "one-api/service"
  13. "one-api/setting/operation_setting"
  14. "sync"
  15. "time"
  16. "github.com/bytedance/gopkg/util/gopool"
  17. "github.com/gin-gonic/gin"
  18. "github.com/gorilla/websocket"
  19. )
  20. func SetupApiRequestHeader(info *common.RelayInfo, c *gin.Context, req *http.Header) {
  21. if info.RelayMode == constant.RelayModeAudioTranscription || info.RelayMode == constant.RelayModeAudioTranslation {
  22. // multipart/form-data
  23. } else if info.RelayMode == constant.RelayModeRealtime {
  24. // websocket
  25. } else {
  26. req.Set("Content-Type", c.Request.Header.Get("Content-Type"))
  27. req.Set("Accept", c.Request.Header.Get("Accept"))
  28. if info.IsStream && c.Request.Header.Get("Accept") == "" {
  29. req.Set("Accept", "text/event-stream")
  30. }
  31. }
  32. }
  33. func DoApiRequest(a Adaptor, c *gin.Context, info *common.RelayInfo, requestBody io.Reader) (*http.Response, error) {
  34. fullRequestURL, err := a.GetRequestURL(info)
  35. if err != nil {
  36. return nil, fmt.Errorf("get request url failed: %w", err)
  37. }
  38. if common2.DebugEnabled {
  39. println("fullRequestURL:", fullRequestURL)
  40. }
  41. req, err := http.NewRequest(c.Request.Method, fullRequestURL, requestBody)
  42. if err != nil {
  43. return nil, fmt.Errorf("new request failed: %w", err)
  44. }
  45. err = a.SetupRequestHeader(c, &req.Header, info)
  46. if err != nil {
  47. return nil, fmt.Errorf("setup request header failed: %w", err)
  48. }
  49. resp, err := doRequest(c, req, info)
  50. if err != nil {
  51. return nil, fmt.Errorf("do request failed: %w", err)
  52. }
  53. return resp, nil
  54. }
  55. func DoFormRequest(a Adaptor, c *gin.Context, info *common.RelayInfo, requestBody io.Reader) (*http.Response, error) {
  56. fullRequestURL, err := a.GetRequestURL(info)
  57. if err != nil {
  58. return nil, fmt.Errorf("get request url failed: %w", err)
  59. }
  60. req, err := http.NewRequest(c.Request.Method, fullRequestURL, requestBody)
  61. if err != nil {
  62. return nil, fmt.Errorf("new request failed: %w", err)
  63. }
  64. // set form data
  65. req.Header.Set("Content-Type", c.Request.Header.Get("Content-Type"))
  66. err = a.SetupRequestHeader(c, &req.Header, info)
  67. if err != nil {
  68. return nil, fmt.Errorf("setup request header failed: %w", err)
  69. }
  70. resp, err := doRequest(c, req, info)
  71. if err != nil {
  72. return nil, fmt.Errorf("do request failed: %w", err)
  73. }
  74. return resp, nil
  75. }
  76. func DoWssRequest(a Adaptor, c *gin.Context, info *common.RelayInfo, requestBody io.Reader) (*websocket.Conn, error) {
  77. fullRequestURL, err := a.GetRequestURL(info)
  78. if err != nil {
  79. return nil, fmt.Errorf("get request url failed: %w", err)
  80. }
  81. targetHeader := http.Header{}
  82. err = a.SetupRequestHeader(c, &targetHeader, info)
  83. if err != nil {
  84. return nil, fmt.Errorf("setup request header failed: %w", err)
  85. }
  86. targetHeader.Set("Content-Type", c.Request.Header.Get("Content-Type"))
  87. targetConn, _, err := websocket.DefaultDialer.Dial(fullRequestURL, targetHeader)
  88. if err != nil {
  89. return nil, fmt.Errorf("dial failed to %s: %w", fullRequestURL, err)
  90. }
  91. // send request body
  92. //all, err := io.ReadAll(requestBody)
  93. //err = service.WssString(c, targetConn, string(all))
  94. return targetConn, nil
  95. }
  96. func doRequest(c *gin.Context, req *http.Request, info *common.RelayInfo) (*http.Response, error) {
  97. var client *http.Client
  98. var err error
  99. if proxyURL, ok := info.ChannelSetting["proxy"]; ok {
  100. client, err = service.NewProxyHttpClient(proxyURL.(string))
  101. if err != nil {
  102. return nil, fmt.Errorf("new proxy http client failed: %w", err)
  103. }
  104. } else {
  105. client = service.GetHttpClient()
  106. }
  107. // 流式请求 ping 保活
  108. var stopPinger func()
  109. generalSettings := operation_setting.GetGeneralSetting()
  110. pingEnabled := generalSettings.PingIntervalEnabled
  111. var pingerWg sync.WaitGroup
  112. if info.IsStream {
  113. helper.SetEventStreamHeaders(c)
  114. pingInterval := time.Duration(generalSettings.PingIntervalSeconds) * time.Second
  115. var pingerCtx context.Context
  116. pingerCtx, stopPinger = context.WithCancel(c.Request.Context())
  117. if pingEnabled {
  118. pingerWg.Add(1)
  119. gopool.Go(func() {
  120. defer pingerWg.Done()
  121. if pingInterval <= 0 {
  122. pingInterval = helper.DefaultPingInterval
  123. }
  124. ticker := time.NewTicker(pingInterval)
  125. defer ticker.Stop()
  126. var pingMutex sync.Mutex
  127. if common2.DebugEnabled {
  128. println("SSE ping goroutine started")
  129. }
  130. for {
  131. select {
  132. case <-ticker.C:
  133. pingMutex.Lock()
  134. err2 := helper.PingData(c)
  135. pingMutex.Unlock()
  136. if err2 != nil {
  137. common2.LogError(c, "SSE ping error: "+err.Error())
  138. return
  139. }
  140. if common2.DebugEnabled {
  141. println("SSE ping data sent.")
  142. }
  143. case <-pingerCtx.Done():
  144. if common2.DebugEnabled {
  145. println("SSE ping goroutine stopped.")
  146. }
  147. return
  148. }
  149. }
  150. })
  151. }
  152. }
  153. resp, err := client.Do(req)
  154. // request结束后停止ping
  155. if info.IsStream && pingEnabled {
  156. stopPinger()
  157. pingerWg.Wait()
  158. }
  159. if err != nil {
  160. return nil, err
  161. }
  162. if resp == nil {
  163. return nil, errors.New("resp is nil")
  164. }
  165. _ = req.Body.Close()
  166. _ = c.Request.Body.Close()
  167. return resp, nil
  168. }
  169. func DoTaskApiRequest(a TaskAdaptor, c *gin.Context, info *common.TaskRelayInfo, requestBody io.Reader) (*http.Response, error) {
  170. fullRequestURL, err := a.BuildRequestURL(info)
  171. if err != nil {
  172. return nil, err
  173. }
  174. req, err := http.NewRequest(c.Request.Method, fullRequestURL, requestBody)
  175. if err != nil {
  176. return nil, fmt.Errorf("new request failed: %w", err)
  177. }
  178. req.GetBody = func() (io.ReadCloser, error) {
  179. return io.NopCloser(requestBody), nil
  180. }
  181. err = a.BuildRequestHeader(c, req, info)
  182. if err != nil {
  183. return nil, fmt.Errorf("setup request header failed: %w", err)
  184. }
  185. resp, err := doRequest(c, req, info.RelayInfo)
  186. if err != nil {
  187. return nil, fmt.Errorf("do request failed: %w", err)
  188. }
  189. return resp, nil
  190. }