proxy_wrapper.go 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250
  1. package proxy
  2. import (
  3. "context"
  4. "fmt"
  5. "net"
  6. "sync"
  7. "sync/atomic"
  8. "time"
  9. "github.com/fatedier/frp/client/event"
  10. "github.com/fatedier/frp/client/health"
  11. "github.com/fatedier/frp/pkg/config"
  12. "github.com/fatedier/frp/pkg/msg"
  13. "github.com/fatedier/frp/pkg/util/xlog"
  14. "github.com/fatedier/golib/errors"
  15. )
  16. const (
  17. ProxyPhaseNew = "new"
  18. ProxyPhaseWaitStart = "wait start"
  19. ProxyPhaseStartErr = "start error"
  20. ProxyPhaseRunning = "running"
  21. ProxyPhaseCheckFailed = "check failed"
  22. ProxyPhaseClosed = "closed"
  23. )
  24. var (
  25. statusCheckInterval time.Duration = 3 * time.Second
  26. waitResponseTimeout = 20 * time.Second
  27. startErrTimeout = 30 * time.Second
  28. )
  29. type WorkingStatus struct {
  30. Name string `json:"name"`
  31. Type string `json:"type"`
  32. Phase string `json:"status"`
  33. Err string `json:"err"`
  34. Cfg config.ProxyConf `json:"cfg"`
  35. // Got from server.
  36. RemoteAddr string `json:"remote_addr"`
  37. }
  38. type Wrapper struct {
  39. WorkingStatus
  40. // underlying proxy
  41. pxy Proxy
  42. // if ProxyConf has healcheck config
  43. // monitor will watch if it is alive
  44. monitor *health.Monitor
  45. // event handler
  46. handler event.Handler
  47. health uint32
  48. lastSendStartMsg time.Time
  49. lastStartErr time.Time
  50. closeCh chan struct{}
  51. healthNotifyCh chan struct{}
  52. mu sync.RWMutex
  53. xl *xlog.Logger
  54. ctx context.Context
  55. }
  56. func NewWrapper(ctx context.Context, cfg config.ProxyConf, clientCfg config.ClientCommonConf, eventHandler event.Handler, serverUDPPort int) *Wrapper {
  57. baseInfo := cfg.GetBaseInfo()
  58. xl := xlog.FromContextSafe(ctx).Spawn().AppendPrefix(baseInfo.ProxyName)
  59. pw := &Wrapper{
  60. WorkingStatus: WorkingStatus{
  61. Name: baseInfo.ProxyName,
  62. Type: baseInfo.ProxyType,
  63. Phase: ProxyPhaseNew,
  64. Cfg: cfg,
  65. },
  66. closeCh: make(chan struct{}),
  67. healthNotifyCh: make(chan struct{}),
  68. handler: eventHandler,
  69. xl: xl,
  70. ctx: xlog.NewContext(ctx, xl),
  71. }
  72. if baseInfo.HealthCheckType != "" {
  73. pw.health = 1 // means failed
  74. pw.monitor = health.NewMonitor(pw.ctx, baseInfo.HealthCheckType, baseInfo.HealthCheckIntervalS,
  75. baseInfo.HealthCheckTimeoutS, baseInfo.HealthCheckMaxFailed, baseInfo.HealthCheckAddr,
  76. baseInfo.HealthCheckURL, pw.statusNormalCallback, pw.statusFailedCallback)
  77. xl.Trace("enable health check monitor")
  78. }
  79. pw.pxy = NewProxy(pw.ctx, pw.Cfg, clientCfg, serverUDPPort)
  80. return pw
  81. }
  82. func (pw *Wrapper) SetRunningStatus(remoteAddr string, respErr string) error {
  83. pw.mu.Lock()
  84. defer pw.mu.Unlock()
  85. if pw.Phase != ProxyPhaseWaitStart {
  86. return fmt.Errorf("status not wait start, ignore start message")
  87. }
  88. pw.RemoteAddr = remoteAddr
  89. if respErr != "" {
  90. pw.Phase = ProxyPhaseStartErr
  91. pw.Err = respErr
  92. pw.lastStartErr = time.Now()
  93. return fmt.Errorf(pw.Err)
  94. }
  95. if err := pw.pxy.Run(); err != nil {
  96. pw.close()
  97. pw.Phase = ProxyPhaseStartErr
  98. pw.Err = err.Error()
  99. pw.lastStartErr = time.Now()
  100. return err
  101. }
  102. pw.Phase = ProxyPhaseRunning
  103. pw.Err = ""
  104. return nil
  105. }
  106. func (pw *Wrapper) Start() {
  107. go pw.checkWorker()
  108. if pw.monitor != nil {
  109. go pw.monitor.Start()
  110. }
  111. }
  112. func (pw *Wrapper) Stop() {
  113. pw.mu.Lock()
  114. defer pw.mu.Unlock()
  115. close(pw.closeCh)
  116. close(pw.healthNotifyCh)
  117. pw.pxy.Close()
  118. if pw.monitor != nil {
  119. pw.monitor.Stop()
  120. }
  121. pw.Phase = ProxyPhaseClosed
  122. pw.close()
  123. }
  124. func (pw *Wrapper) close() {
  125. pw.handler(event.EvCloseProxy, &event.CloseProxyPayload{
  126. CloseProxyMsg: &msg.CloseProxy{
  127. ProxyName: pw.Name,
  128. },
  129. })
  130. }
  131. func (pw *Wrapper) checkWorker() {
  132. xl := pw.xl
  133. if pw.monitor != nil {
  134. // let monitor do check request first
  135. time.Sleep(500 * time.Millisecond)
  136. }
  137. for {
  138. // check proxy status
  139. now := time.Now()
  140. if atomic.LoadUint32(&pw.health) == 0 {
  141. pw.mu.Lock()
  142. if pw.Phase == ProxyPhaseNew ||
  143. pw.Phase == ProxyPhaseCheckFailed ||
  144. (pw.Phase == ProxyPhaseWaitStart && now.After(pw.lastSendStartMsg.Add(waitResponseTimeout))) ||
  145. (pw.Phase == ProxyPhaseStartErr && now.After(pw.lastStartErr.Add(startErrTimeout))) {
  146. xl.Trace("change status from [%s] to [%s]", pw.Phase, ProxyPhaseWaitStart)
  147. pw.Phase = ProxyPhaseWaitStart
  148. var newProxyMsg msg.NewProxy
  149. pw.Cfg.MarshalToMsg(&newProxyMsg)
  150. pw.lastSendStartMsg = now
  151. pw.handler(event.EvStartProxy, &event.StartProxyPayload{
  152. NewProxyMsg: &newProxyMsg,
  153. })
  154. }
  155. pw.mu.Unlock()
  156. } else {
  157. pw.mu.Lock()
  158. if pw.Phase == ProxyPhaseRunning || pw.Phase == ProxyPhaseWaitStart {
  159. pw.close()
  160. xl.Trace("change status from [%s] to [%s]", pw.Phase, ProxyPhaseCheckFailed)
  161. pw.Phase = ProxyPhaseCheckFailed
  162. }
  163. pw.mu.Unlock()
  164. }
  165. select {
  166. case <-pw.closeCh:
  167. return
  168. case <-time.After(statusCheckInterval):
  169. case <-pw.healthNotifyCh:
  170. }
  171. }
  172. }
  173. func (pw *Wrapper) statusNormalCallback() {
  174. xl := pw.xl
  175. atomic.StoreUint32(&pw.health, 0)
  176. errors.PanicToError(func() {
  177. select {
  178. case pw.healthNotifyCh <- struct{}{}:
  179. default:
  180. }
  181. })
  182. xl.Info("health check success")
  183. }
  184. func (pw *Wrapper) statusFailedCallback() {
  185. xl := pw.xl
  186. atomic.StoreUint32(&pw.health, 1)
  187. errors.PanicToError(func() {
  188. select {
  189. case pw.healthNotifyCh <- struct{}{}:
  190. default:
  191. }
  192. })
  193. xl.Info("health check failed")
  194. }
  195. func (pw *Wrapper) InWorkConn(workConn net.Conn, m *msg.StartWorkConn) {
  196. xl := pw.xl
  197. pw.mu.RLock()
  198. pxy := pw.pxy
  199. pw.mu.RUnlock()
  200. if pxy != nil && pw.Phase == ProxyPhaseRunning {
  201. xl.Debug("start a new work connection, localAddr: %s remoteAddr: %s", workConn.LocalAddr().String(), workConn.RemoteAddr().String())
  202. go pxy.InWorkConn(workConn, m)
  203. } else {
  204. workConn.Close()
  205. }
  206. }
  207. func (pw *Wrapper) GetStatus() *WorkingStatus {
  208. pw.mu.RLock()
  209. defer pw.mu.RUnlock()
  210. ps := &WorkingStatus{
  211. Name: pw.Name,
  212. Type: pw.Type,
  213. Phase: pw.Phase,
  214. Err: pw.Err,
  215. Cfg: pw.Cfg,
  216. RemoteAddr: pw.RemoteAddr,
  217. }
  218. return ps
  219. }