visitor.go 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553
  1. // Copyright 2017 fatedier, fatedier@gmail.com
  2. //
  3. // Licensed under the Apache License, Version 2.0 (the "License");
  4. // you may not use this file except in compliance with the License.
  5. // You may obtain a copy of the License at
  6. //
  7. // http://www.apache.org/licenses/LICENSE-2.0
  8. //
  9. // Unless required by applicable law or agreed to in writing, software
  10. // distributed under the License is distributed on an "AS IS" BASIS,
  11. // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  12. // See the License for the specific language governing permissions and
  13. // limitations under the License.
  14. package client
  15. import (
  16. "bytes"
  17. "context"
  18. "fmt"
  19. "io"
  20. "net"
  21. "sync"
  22. "time"
  23. "github.com/fatedier/frp/pkg/config"
  24. "github.com/fatedier/frp/pkg/msg"
  25. "github.com/fatedier/frp/pkg/proto/udp"
  26. frpNet "github.com/fatedier/frp/pkg/util/net"
  27. "github.com/fatedier/frp/pkg/util/util"
  28. "github.com/fatedier/frp/pkg/util/xlog"
  29. "github.com/fatedier/golib/errors"
  30. frpIo "github.com/fatedier/golib/io"
  31. "github.com/fatedier/golib/pool"
  32. fmux "github.com/hashicorp/yamux"
  33. )
  34. // Visitor is used for forward traffics from local port tot remote service.
  35. type Visitor interface {
  36. Run() error
  37. Close()
  38. }
  39. func NewVisitor(ctx context.Context, ctl *Control, cfg config.VisitorConf) (visitor Visitor) {
  40. xl := xlog.FromContextSafe(ctx).Spawn().AppendPrefix(cfg.GetBaseInfo().ProxyName)
  41. baseVisitor := BaseVisitor{
  42. ctl: ctl,
  43. ctx: xlog.NewContext(ctx, xl),
  44. }
  45. switch cfg := cfg.(type) {
  46. case *config.STCPVisitorConf:
  47. visitor = &STCPVisitor{
  48. BaseVisitor: &baseVisitor,
  49. cfg: cfg,
  50. }
  51. case *config.XTCPVisitorConf:
  52. visitor = &XTCPVisitor{
  53. BaseVisitor: &baseVisitor,
  54. cfg: cfg,
  55. }
  56. case *config.SUDPVisitorConf:
  57. visitor = &SUDPVisitor{
  58. BaseVisitor: &baseVisitor,
  59. cfg: cfg,
  60. checkCloseCh: make(chan struct{}),
  61. }
  62. }
  63. return
  64. }
  65. type BaseVisitor struct {
  66. ctl *Control
  67. l net.Listener
  68. closed bool
  69. mu sync.RWMutex
  70. ctx context.Context
  71. }
  72. type STCPVisitor struct {
  73. *BaseVisitor
  74. cfg *config.STCPVisitorConf
  75. }
  76. func (sv *STCPVisitor) Run() (err error) {
  77. sv.l, err = net.Listen("tcp", fmt.Sprintf("%s:%d", sv.cfg.BindAddr, sv.cfg.BindPort))
  78. if err != nil {
  79. return
  80. }
  81. go sv.worker()
  82. return
  83. }
  84. func (sv *STCPVisitor) Close() {
  85. sv.l.Close()
  86. }
  87. func (sv *STCPVisitor) worker() {
  88. xl := xlog.FromContextSafe(sv.ctx)
  89. for {
  90. conn, err := sv.l.Accept()
  91. if err != nil {
  92. xl.Warn("stcp local listener closed")
  93. return
  94. }
  95. go sv.handleConn(conn)
  96. }
  97. }
  98. func (sv *STCPVisitor) handleConn(userConn net.Conn) {
  99. xl := xlog.FromContextSafe(sv.ctx)
  100. defer userConn.Close()
  101. xl.Debug("get a new stcp user connection")
  102. visitorConn, err := sv.ctl.connectServer()
  103. if err != nil {
  104. return
  105. }
  106. defer visitorConn.Close()
  107. now := time.Now().Unix()
  108. newVisitorConnMsg := &msg.NewVisitorConn{
  109. ProxyName: sv.cfg.ServerName,
  110. SignKey: util.GetAuthKey(sv.cfg.Sk, now),
  111. Timestamp: now,
  112. UseEncryption: sv.cfg.UseEncryption,
  113. UseCompression: sv.cfg.UseCompression,
  114. }
  115. err = msg.WriteMsg(visitorConn, newVisitorConnMsg)
  116. if err != nil {
  117. xl.Warn("send newVisitorConnMsg to server error: %v", err)
  118. return
  119. }
  120. var newVisitorConnRespMsg msg.NewVisitorConnResp
  121. visitorConn.SetReadDeadline(time.Now().Add(10 * time.Second))
  122. err = msg.ReadMsgInto(visitorConn, &newVisitorConnRespMsg)
  123. if err != nil {
  124. xl.Warn("get newVisitorConnRespMsg error: %v", err)
  125. return
  126. }
  127. visitorConn.SetReadDeadline(time.Time{})
  128. if newVisitorConnRespMsg.Error != "" {
  129. xl.Warn("start new visitor connection error: %s", newVisitorConnRespMsg.Error)
  130. return
  131. }
  132. var remote io.ReadWriteCloser
  133. remote = visitorConn
  134. if sv.cfg.UseEncryption {
  135. remote, err = frpIo.WithEncryption(remote, []byte(sv.cfg.Sk))
  136. if err != nil {
  137. xl.Error("create encryption stream error: %v", err)
  138. return
  139. }
  140. }
  141. if sv.cfg.UseCompression {
  142. remote = frpIo.WithCompression(remote)
  143. }
  144. frpIo.Join(userConn, remote)
  145. }
  146. type XTCPVisitor struct {
  147. *BaseVisitor
  148. cfg *config.XTCPVisitorConf
  149. }
  150. func (sv *XTCPVisitor) Run() (err error) {
  151. sv.l, err = net.Listen("tcp", fmt.Sprintf("%s:%d", sv.cfg.BindAddr, sv.cfg.BindPort))
  152. if err != nil {
  153. return
  154. }
  155. go sv.worker()
  156. return
  157. }
  158. func (sv *XTCPVisitor) Close() {
  159. sv.l.Close()
  160. }
  161. func (sv *XTCPVisitor) worker() {
  162. xl := xlog.FromContextSafe(sv.ctx)
  163. for {
  164. conn, err := sv.l.Accept()
  165. if err != nil {
  166. xl.Warn("xtcp local listener closed")
  167. return
  168. }
  169. go sv.handleConn(conn)
  170. }
  171. }
  172. func (sv *XTCPVisitor) handleConn(userConn net.Conn) {
  173. xl := xlog.FromContextSafe(sv.ctx)
  174. defer userConn.Close()
  175. xl.Debug("get a new xtcp user connection")
  176. if sv.ctl.serverUDPPort == 0 {
  177. xl.Error("xtcp is not supported by server")
  178. return
  179. }
  180. raddr, err := net.ResolveUDPAddr("udp",
  181. fmt.Sprintf("%s:%d", sv.ctl.clientCfg.ServerAddr, sv.ctl.serverUDPPort))
  182. if err != nil {
  183. xl.Error("resolve server UDP addr error")
  184. return
  185. }
  186. visitorConn, err := net.DialUDP("udp", nil, raddr)
  187. if err != nil {
  188. xl.Warn("dial server udp addr error: %v", err)
  189. return
  190. }
  191. defer visitorConn.Close()
  192. now := time.Now().Unix()
  193. natHoleVisitorMsg := &msg.NatHoleVisitor{
  194. ProxyName: sv.cfg.ServerName,
  195. SignKey: util.GetAuthKey(sv.cfg.Sk, now),
  196. Timestamp: now,
  197. }
  198. err = msg.WriteMsg(visitorConn, natHoleVisitorMsg)
  199. if err != nil {
  200. xl.Warn("send natHoleVisitorMsg to server error: %v", err)
  201. return
  202. }
  203. // Wait for client address at most 10 seconds.
  204. var natHoleRespMsg msg.NatHoleResp
  205. visitorConn.SetReadDeadline(time.Now().Add(10 * time.Second))
  206. buf := pool.GetBuf(1024)
  207. n, err := visitorConn.Read(buf)
  208. if err != nil {
  209. xl.Warn("get natHoleRespMsg error: %v", err)
  210. return
  211. }
  212. err = msg.ReadMsgInto(bytes.NewReader(buf[:n]), &natHoleRespMsg)
  213. if err != nil {
  214. xl.Warn("get natHoleRespMsg error: %v", err)
  215. return
  216. }
  217. visitorConn.SetReadDeadline(time.Time{})
  218. pool.PutBuf(buf)
  219. if natHoleRespMsg.Error != "" {
  220. xl.Error("natHoleRespMsg get error info: %s", natHoleRespMsg.Error)
  221. return
  222. }
  223. xl.Trace("get natHoleRespMsg, sid [%s], client address [%s], visitor address [%s]", natHoleRespMsg.Sid, natHoleRespMsg.ClientAddr, natHoleRespMsg.VisitorAddr)
  224. // Close visitorConn, so we can use it's local address.
  225. visitorConn.Close()
  226. // send sid message to client
  227. laddr, _ := net.ResolveUDPAddr("udp", visitorConn.LocalAddr().String())
  228. daddr, err := net.ResolveUDPAddr("udp", natHoleRespMsg.ClientAddr)
  229. if err != nil {
  230. xl.Error("resolve client udp address error: %v", err)
  231. return
  232. }
  233. lConn, err := net.DialUDP("udp", laddr, daddr)
  234. if err != nil {
  235. xl.Error("dial client udp address error: %v", err)
  236. return
  237. }
  238. defer lConn.Close()
  239. lConn.Write([]byte(natHoleRespMsg.Sid))
  240. // read ack sid from client
  241. sidBuf := pool.GetBuf(1024)
  242. lConn.SetReadDeadline(time.Now().Add(8 * time.Second))
  243. n, err = lConn.Read(sidBuf)
  244. if err != nil {
  245. xl.Warn("get sid from client error: %v", err)
  246. return
  247. }
  248. lConn.SetReadDeadline(time.Time{})
  249. if string(sidBuf[:n]) != natHoleRespMsg.Sid {
  250. xl.Warn("incorrect sid from client")
  251. return
  252. }
  253. pool.PutBuf(sidBuf)
  254. xl.Info("nat hole connection make success, sid [%s]", natHoleRespMsg.Sid)
  255. // wrap kcp connection
  256. var remote io.ReadWriteCloser
  257. remote, err = frpNet.NewKCPConnFromUDP(lConn, true, natHoleRespMsg.ClientAddr)
  258. if err != nil {
  259. xl.Error("create kcp connection from udp connection error: %v", err)
  260. return
  261. }
  262. fmuxCfg := fmux.DefaultConfig()
  263. fmuxCfg.KeepAliveInterval = 5 * time.Second
  264. fmuxCfg.LogOutput = io.Discard
  265. sess, err := fmux.Client(remote, fmuxCfg)
  266. if err != nil {
  267. xl.Error("create yamux session error: %v", err)
  268. return
  269. }
  270. defer sess.Close()
  271. muxConn, err := sess.Open()
  272. if err != nil {
  273. xl.Error("open yamux stream error: %v", err)
  274. return
  275. }
  276. var muxConnRWCloser io.ReadWriteCloser = muxConn
  277. if sv.cfg.UseEncryption {
  278. muxConnRWCloser, err = frpIo.WithEncryption(muxConnRWCloser, []byte(sv.cfg.Sk))
  279. if err != nil {
  280. xl.Error("create encryption stream error: %v", err)
  281. return
  282. }
  283. }
  284. if sv.cfg.UseCompression {
  285. muxConnRWCloser = frpIo.WithCompression(muxConnRWCloser)
  286. }
  287. frpIo.Join(userConn, muxConnRWCloser)
  288. xl.Debug("join connections closed")
  289. }
  290. type SUDPVisitor struct {
  291. *BaseVisitor
  292. checkCloseCh chan struct{}
  293. // udpConn is the listener of udp packet
  294. udpConn *net.UDPConn
  295. readCh chan *msg.UDPPacket
  296. sendCh chan *msg.UDPPacket
  297. cfg *config.SUDPVisitorConf
  298. }
  299. // SUDP Run start listen a udp port
  300. func (sv *SUDPVisitor) Run() (err error) {
  301. xl := xlog.FromContextSafe(sv.ctx)
  302. addr, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", sv.cfg.BindAddr, sv.cfg.BindPort))
  303. if err != nil {
  304. return fmt.Errorf("sudp ResolveUDPAddr error: %v", err)
  305. }
  306. sv.udpConn, err = net.ListenUDP("udp", addr)
  307. if err != nil {
  308. return fmt.Errorf("listen udp port %s error: %v", addr.String(), err)
  309. }
  310. sv.sendCh = make(chan *msg.UDPPacket, 1024)
  311. sv.readCh = make(chan *msg.UDPPacket, 1024)
  312. xl.Info("sudp start to work, listen on %s", addr)
  313. go sv.dispatcher()
  314. go udp.ForwardUserConn(sv.udpConn, sv.readCh, sv.sendCh, int(sv.ctl.clientCfg.UDPPacketSize))
  315. return
  316. }
  317. func (sv *SUDPVisitor) dispatcher() {
  318. xl := xlog.FromContextSafe(sv.ctx)
  319. for {
  320. // loop for get frpc to frps tcp conn
  321. // setup worker
  322. // wait worker to finished
  323. // retry or exit
  324. visitorConn, err := sv.getNewVisitorConn()
  325. if err != nil {
  326. // check if proxy is closed
  327. // if checkCloseCh is close, we will return, other case we will continue to reconnect
  328. select {
  329. case <-sv.checkCloseCh:
  330. xl.Info("frpc sudp visitor proxy is closed")
  331. return
  332. default:
  333. }
  334. time.Sleep(3 * time.Second)
  335. xl.Warn("newVisitorConn to frps error: %v, try to reconnect", err)
  336. continue
  337. }
  338. sv.worker(visitorConn)
  339. select {
  340. case <-sv.checkCloseCh:
  341. return
  342. default:
  343. }
  344. }
  345. }
  346. func (sv *SUDPVisitor) worker(workConn net.Conn) {
  347. xl := xlog.FromContextSafe(sv.ctx)
  348. xl.Debug("starting sudp proxy worker")
  349. wg := &sync.WaitGroup{}
  350. wg.Add(2)
  351. closeCh := make(chan struct{})
  352. // udp service -> frpc -> frps -> frpc visitor -> user
  353. workConnReaderFn := func(conn net.Conn) {
  354. defer func() {
  355. conn.Close()
  356. close(closeCh)
  357. wg.Done()
  358. }()
  359. for {
  360. var (
  361. rawMsg msg.Message
  362. errRet error
  363. )
  364. // frpc will send heartbeat in workConn to frpc visitor for keeping alive
  365. conn.SetReadDeadline(time.Now().Add(60 * time.Second))
  366. if rawMsg, errRet = msg.ReadMsg(conn); errRet != nil {
  367. xl.Warn("read from workconn for user udp conn error: %v", errRet)
  368. return
  369. }
  370. conn.SetReadDeadline(time.Time{})
  371. switch m := rawMsg.(type) {
  372. case *msg.Ping:
  373. xl.Debug("frpc visitor get ping message from frpc")
  374. continue
  375. case *msg.UDPPacket:
  376. if errRet := errors.PanicToError(func() {
  377. sv.readCh <- m
  378. xl.Trace("frpc visitor get udp packet from workConn: %s", m.Content)
  379. }); errRet != nil {
  380. xl.Info("reader goroutine for udp work connection closed")
  381. return
  382. }
  383. }
  384. }
  385. }
  386. // udp service <- frpc <- frps <- frpc visitor <- user
  387. workConnSenderFn := func(conn net.Conn) {
  388. defer func() {
  389. conn.Close()
  390. wg.Done()
  391. }()
  392. var errRet error
  393. for {
  394. select {
  395. case udpMsg, ok := <-sv.sendCh:
  396. if !ok {
  397. xl.Info("sender goroutine for udp work connection closed")
  398. return
  399. }
  400. if errRet = msg.WriteMsg(conn, udpMsg); errRet != nil {
  401. xl.Warn("sender goroutine for udp work connection closed: %v", errRet)
  402. return
  403. }
  404. xl.Trace("send udp package to workConn: %s", udpMsg.Content)
  405. case <-closeCh:
  406. return
  407. }
  408. }
  409. }
  410. go workConnReaderFn(workConn)
  411. go workConnSenderFn(workConn)
  412. wg.Wait()
  413. xl.Info("sudp worker is closed")
  414. }
  415. func (sv *SUDPVisitor) getNewVisitorConn() (net.Conn, error) {
  416. xl := xlog.FromContextSafe(sv.ctx)
  417. visitorConn, err := sv.ctl.connectServer()
  418. if err != nil {
  419. return nil, fmt.Errorf("frpc connect frps error: %v", err)
  420. }
  421. now := time.Now().Unix()
  422. newVisitorConnMsg := &msg.NewVisitorConn{
  423. ProxyName: sv.cfg.ServerName,
  424. SignKey: util.GetAuthKey(sv.cfg.Sk, now),
  425. Timestamp: now,
  426. UseEncryption: sv.cfg.UseEncryption,
  427. UseCompression: sv.cfg.UseCompression,
  428. }
  429. err = msg.WriteMsg(visitorConn, newVisitorConnMsg)
  430. if err != nil {
  431. return nil, fmt.Errorf("frpc send newVisitorConnMsg to frps error: %v", err)
  432. }
  433. var newVisitorConnRespMsg msg.NewVisitorConnResp
  434. visitorConn.SetReadDeadline(time.Now().Add(10 * time.Second))
  435. err = msg.ReadMsgInto(visitorConn, &newVisitorConnRespMsg)
  436. if err != nil {
  437. return nil, fmt.Errorf("frpc read newVisitorConnRespMsg error: %v", err)
  438. }
  439. visitorConn.SetReadDeadline(time.Time{})
  440. if newVisitorConnRespMsg.Error != "" {
  441. return nil, fmt.Errorf("start new visitor connection error: %s", newVisitorConnRespMsg.Error)
  442. }
  443. var remote io.ReadWriteCloser
  444. remote = visitorConn
  445. if sv.cfg.UseEncryption {
  446. remote, err = frpIo.WithEncryption(remote, []byte(sv.cfg.Sk))
  447. if err != nil {
  448. xl.Error("create encryption stream error: %v", err)
  449. return nil, err
  450. }
  451. }
  452. if sv.cfg.UseCompression {
  453. remote = frpIo.WithCompression(remote)
  454. }
  455. return frpNet.WrapReadWriteCloserToConn(remote, visitorConn), nil
  456. }
  457. func (sv *SUDPVisitor) Close() {
  458. sv.mu.Lock()
  459. defer sv.mu.Unlock()
  460. select {
  461. case <-sv.checkCloseCh:
  462. return
  463. default:
  464. close(sv.checkCloseCh)
  465. }
  466. if sv.udpConn != nil {
  467. sv.udpConn.Close()
  468. }
  469. close(sv.readCh)
  470. close(sv.sendCh)
  471. }