Skip to content
File

Blob: tun/client/client.go

go174 lines
1package client
2 
3import (
4 "context"
5 "net"
6 "os"
7 "sync"
8 "time"
9 
10 "go.miragespace.co/specter/overlay"
11 "go.miragespace.co/specter/spec/protocol"
12 "go.miragespace.co/specter/spec/rpc"
13 "go.miragespace.co/specter/spec/rtt"
14 "go.miragespace.co/specter/spec/transport"
15 "go.miragespace.co/specter/spec/tun"
16 "go.miragespace.co/specter/util"
17 "go.miragespace.co/specter/util/acceptor"
18 
19 "github.com/Yiling-J/theine-go"
20 "github.com/zhangyunhao116/skipmap"
21 "go.uber.org/atomic"
22 "go.uber.org/zap"
23)
24 
25var (
26 checkInterval = time.Second * 30
27 rttInterval = transport.RTTMeasureInterval
28 certCheckInterval = time.Hour * 24 // Check certificate daily for long-running clients
29)
30 
31const (
32 connectTimeout = time.Second * 5
33 rpcTimeout = time.Second * 5
34 renewalWindow = 30 * 24 * time.Hour // Renew certificates 30 days before expiry
35)
36 
37type KeylessProxyConfig struct {
38 HTTPListner net.Listener
39 HTTPSListner net.Listener
40 ALPNMux *overlay.ALPNMux
41}
42 
43type ClientConfig struct {
44 Logger *zap.Logger
45 Configuration *Config
46 PKIClient protocol.PKIService
47 ServerTransport transport.Transport
48 Recorder rtt.Recorder
49 ReloadSignal <-chan os.Signal
50 ServerListener net.Listener
51 KeylessProxy KeylessProxyConfig
52}
53 
54type Client struct {
55 ClientConfig
56 *forwarder
57 configMu sync.RWMutex
58 closeWg sync.WaitGroup
59 syncMu sync.Mutex
60 syncStateMu sync.RWMutex
61 publication map[string]publicationState
62 hostnameNextCandidate uint64
63 lastSync SyncResult
64 nextSync time.Time
65 syncBackoff time.Duration
66 tunnelClient rpc.TunnelClient
67 parentCtx context.Context
68 connections *skipmap.StringMap[*protocol.Node]
69 rpcAcceptor *acceptor.HTTP2Acceptor
70 keylessCertificateCache *theine.LoadingCache[string, keylessCertificateResult]
71 closeCh chan struct{}
72 closed atomic.Bool
73}
74 
75func NewClient(ctx context.Context, cfg ClientConfig) (*Client, error) {
76 c := &Client{
77 ClientConfig: cfg,
78 parentCtx: ctx,
79 forwarder: newForwarder(cfg.Logger),
80 connections: skipmap.NewString[*protocol.Node](),
81 tunnelClient: rpc.DynamicTunnelClient(ctx, cfg.ServerTransport),
82 rpcAcceptor: acceptor.NewH2Acceptor(nil),
83 closeCh: make(chan struct{}),
84 }
85 
86 if c.Configuration.Certificate != "" {
87 c.configMu.RLock()
88 apex := c.ClientConfig.Configuration.Apex
89 c.configMu.RUnlock()
90 if err := c.updateTransportCert(); err != nil {
91 return nil, err
92 }
93 c.Logger = c.Logger.With(zap.Uint64("id", c.ServerTransport.Identity().GetId()))
94 c.forwarder.logger = c.Logger
95 if err := c.bootstrap(ctx, apex); err != nil {
96 return nil, err
97 }
98 }
99 
100 keylessCache, err := theine.NewBuilder[string, keylessCertificateResult](cacheTotalCost).
101 BuildWithLoader(c.keylesCertificateCacheLoader)
102 if err != nil {
103 return nil, err
104 }
105 c.keylessCertificateCache = keylessCache
106 
107 return c, nil
108}
109 
110func (c *Client) Initialize(ctx context.Context, syncTunnels bool) error {
111 if err := c.maintainConnections(ctx); err != nil {
112 return err
113 }
114 
115 c.Logger.Info("Waiting for RTT measurement...", zap.Duration("max", rttInterval))
116 time.Sleep(util.RandomTimeRange(rttInterval))
117 
118 if syncTunnels {
119 if result := c.SyncConfigTunnels(ctx); result.Error != "" {
120 c.Logger.Warn("Initial tunnel synchronization is pending retry", zap.String("error", result.Error))
121 }
122 }
123 return nil
124}
125 
126func (c *Client) GetConnectedNodes() []*protocol.Node {
127 return c.getConnectedNodes()
128}
129 
130func (c *Client) Start(ctx context.Context) {
131 c.Logger.Info("Listening for tunnel traffic")
132 
133 c.closeWg.Add(3)
134 
135 streamRouter := transport.NewStreamRouter(c.Logger, nil, c.ServerTransport)
136 streamRouter.HandleTunnel(protocol.Stream_DIRECT, func(delegation *transport.StreamDelegate) {
137 link, ok := receiveLink(c.Logger, delegation)
138 if !ok {
139 return
140 }
141 c.handleIncomingDelegation(ctx, link, delegation)
142 })
143 c.attachRPC(ctx, streamRouter)
144 
145 go streamRouter.Accept(ctx)
146 go c.periodicReconnection(ctx)
147 go c.reloadOnSignal(ctx)
148 go c.certificateMaintainer(ctx)
149 go c.startLocalServer(ctx)
150 go c.startKeylessProxy()
151}
152 
153func (c *Client) handleIncomingDelegation(ctx context.Context, link *protocol.Link, delegation net.Conn) error {
154 hostname := link.GetHostname()
155 u, ok := c.Configuration.router.Load(hostname)
156 if !ok {
157 c.Logger.Error("Unknown hostname in connection", zap.String("hostname", hostname))
158 delegation.Close()
159 return tun.ErrDestinationNotFound
160 }
161 
162 return c.handleLink(ctx, link, delegation, u)
163}
164 
165func (c *Client) Close() {
166 if !c.closed.CompareAndSwap(false, true) {
167 return
168 }
169 c.rpcAcceptor.Close()
170 c.closeAll()
171 close(c.closeCh)
172 c.closeWg.Wait()
173}