Skip to content
File

Blob: tun/client/forwarder.go

go70 lines
1package client
2 
3import (
4 "context"
5 "net"
6 
7 "go.miragespace.co/specter/spec/protocol"
8 "go.miragespace.co/specter/spec/rpc"
9 "go.miragespace.co/specter/spec/tun"
10 
11 "github.com/zhangyunhao116/skipmap"
12 "go.uber.org/atomic"
13 "go.uber.org/zap"
14)
15 
16type forwarder struct {
17 logger *zap.Logger
18 rootDomain *atomic.String
19 proxies *skipmap.StringMap[*httpProxy]
20}
21 
22func newForwarder(logger *zap.Logger) *forwarder {
23 return &forwarder{
24 logger: logger,
25 rootDomain: atomic.NewString(""),
26 proxies: skipmap.NewString[*httpProxy](),
27 }
28}
29 
30func receiveLink(logger *zap.Logger, delegation net.Conn) (*protocol.Link, bool) {
31 link := &protocol.Link{}
32 if err := rpc.BoundedReceive(delegation, link, 1024); err != nil {
33 logger.Error("Receiving link information from gateway", zap.Error(err))
34 delegation.Close()
35 return nil, false
36 }
37 return link, true
38}
39 
40func (f *forwarder) handleLink(ctx context.Context, link *protocol.Link, delegation net.Conn, u route) error {
41 hostname := link.GetHostname()
42 f.logger.Info("Incoming connection from gateway",
43 zap.String("protocol", link.GetAlpn().String()),
44 zap.String("hostname", link.GetHostname()),
45 zap.String("remote", link.GetRemote()))
46 
47 switch link.GetAlpn() {
48 case protocol.Link_HTTP:
49 f.getHTTPProxy(ctx, hostname, u).acceptor.Handle(delegation)
50 
51 case protocol.Link_TCP:
52 f.forwardStream(ctx, hostname, delegation, u)
53 
54 default:
55 f.logger.Error("Unknown alpn for forwarding", zap.String("alpn", link.GetAlpn().String()))
56 delegation.Close()
57 return tun.ErrDestinationNotFound
58 }
59 
60 return nil
61}
62 
63func (f *forwarder) closeAll() {
64 f.proxies.Range(func(key string, proxy *httpProxy) bool {
65 f.logger.Info("Shutting down proxy", zap.String("hostname", key))
66 proxy.acceptor.Close()
67 return true
68 })
69}