Skip to content
File

Blob: cmd/server/server.go

go886 lines
1package server
2 
3import (
4 "context"
5 "crypto/tls"
6 "crypto/x509"
7 "encoding/base64"
8 "fmt"
9 "net"
10 "net/http"
11 "net/url"
12 "os"
13 "os/signal"
14 "path/filepath"
15 "strconv"
16 "syscall"
17 "time"
18 
19 "go.miragespace.co/specter/acme"
20 chordImpl "go.miragespace.co/specter/chord"
21 cmdlisten "go.miragespace.co/specter/cmd/internal/listen"
22 "go.miragespace.co/specter/gateway"
23 "go.miragespace.co/specter/overlay"
24 "go.miragespace.co/specter/pki"
25 "go.miragespace.co/specter/rtt"
26 acmeSpec "go.miragespace.co/specter/spec/acme"
27 "go.miragespace.co/specter/spec/chord"
28 "go.miragespace.co/specter/spec/cipher"
29 "go.miragespace.co/specter/spec/protocol"
30 "go.miragespace.co/specter/spec/rpc"
31 "go.miragespace.co/specter/spec/transport"
32 "go.miragespace.co/specter/spec/transport/q"
33 "go.miragespace.co/specter/spec/tun"
34 "go.miragespace.co/specter/timing"
35 "go.miragespace.co/specter/tun/server"
36 "go.miragespace.co/specter/util"
37 "go.miragespace.co/specter/util/migrator"
38 "go.miragespace.co/specter/util/reuse"
39 
40 "github.com/TheZeroSlave/zapsentry"
41 "github.com/alecthomas/units"
42 "github.com/getsentry/sentry-go"
43 "github.com/pires/go-proxyproto"
44 "github.com/quic-go/quic-go"
45 "github.com/urfave/cli/v3"
46 "go.uber.org/zap"
47 "go.uber.org/zap/zapcore"
48)
49 
50func Generate() *cli.Command {
51 ip := util.GetOutboundIP()
52 return &cli.Command{
53 Name: "server",
54 Usage: "start an specter server on the edge",
55 Description: `Start a Specter server that joins a Chord DHT ring, routes all tunnels over QUIC, issues/renews TLS certificates via ACME DNS-01, and persists state for high availability.
56
57 Specter server provides an internal endpoint on /_internal under apex domain. To enable internal endpoint, provide username and password
58 under environment variables INTERNAL_USER and INTERNAL_PASS. Absent of them will disable the internal endpoint entirely.
59
60 Warning: do not use certificates issued by public CA for inter-node certificates, otherwise anyone can join your specter network`,
61 ArgsUsage: " ",
62 Flags: []cli.Flag{
63 &cli.StringFlag{
64 Name: "advertise-addr",
65 Aliases: []string{"advertise"},
66 Sources: cli.EnvVars("ADVERTISE_ADDR"),
67 DefaultText: "same as listen-addr",
68 Value: fmt.Sprintf("%s:443", ip.String()),
69 Usage: `Address and port to advertise to specter servers and clients to connect to.
70 Note that specter will use advertised address to derive its Identity hash.`,
71 Category: "Network Options",
72 },
73 &cli.StringSliceFlag{
74 Name: "listen-addr",
75 Aliases: []string{"listen"},
76 Value: []string{fmt.Sprintf("%s:443", ip.String())},
77 Usage: `Repeatable address:port to listen for specter server, specter client and gateway connections. Each entry serves both TCP and UDP unless overridden.
78 Note that if specter is listening on port 443, it will also listen on port 80 to handle http connect proxy, and redirect other http requests to https`,
79 Category: "Network Options",
80 Sources: cli.EnvVars("LISTEN_ADDR"),
81 },
82 &cli.StringSliceFlag{
83 Name: "listen-tcp",
84 DefaultText: "same as listen-addr",
85 Usage: "Override the listen address and port for TCP (repeatable)",
86 Category: "Network Options",
87 Sources: cli.EnvVars("LISTEN_TCP"),
88 },
89 &cli.StringSliceFlag{
90 Name: "listen-udp",
91 DefaultText: "same as listen-addr",
92 Usage: "Override the listen address and port for UDP (repeatable). Required if environment needs a specific address, such as on fly.io",
93 Category: "Network Options",
94 Sources: cli.EnvVars("LISTEN_UDP"),
95 },
96 &cli.BoolFlag{
97 Name: "proxy-protocol",
98 Value: false,
99 Usage: "Parse client IP via PROXY protocol (v1 or v2) when handling TCP connections. Required if environment is behind a TCP Load Balancer, such as on fly.io",
100 Category: "Network Options",
101 },
102 
103 &cli.StringFlag{
104 Name: "listen-rpc",
105 Value: "tcp://127.0.0.1:11180",
106 Usage: `Expose chord's RPC for VNode and KV to an external program. This is required to use with specter's acme dns.
107 NOTE: The listener is exposed without any authentication or authorization. You should only expose it to localhost or unix socket`,
108 Category: "Server Options",
109 },
110 &cli.StringFlag{
111 Name: "data-dir",
112 Aliases: []string{"data"},
113 Usage: "Path to directory that will be used for persisting non-volatile KV data",
114 Required: true,
115 Category: "Server Options",
116 },
117 &cli.StringFlag{
118 Name: "cert-dir",
119 Aliases: []string{"cert"},
120 Usage: `Path to directory containing ca.crt, client-ca.crt, client-ca.key, node.crt, and node.key for mutual TLS between specter server nodes`,
121 Category: "Server Options",
122 },
123 &cli.BoolFlag{
124 Name: "cert-env",
125 Usage: `Load ca.crt (CERT_CA), client-ca.crt (CERT_CLIENT_CA), client-ca.key (CERT_CLIENT_CA_KEY), node.crt (CERT_NODE), and node.key (CERT_NODE_KEY) from environment variables encoded as base64.
126 This can be set instead of loading from cert-dir. Required if environment prefers loading secrets from ENV, such as on fly.io`,
127 Category: "Server Options",
128 },
129 &cli.StringFlag{
130 Name: "sentry",
131 DefaultText: "https://public@sentry.example.com/1",
132 Usage: "Sentry DSN for error monitoring. Alternatively, you can set the DSN via the environment variable SENTRY_DSN",
133 Sources: cli.EnvVars("SENTRY_DSN"),
134 Category: "Server Options",
135 },
136 
137 &cli.IntFlag{
138 Name: "virtual",
139 Usage: "Number of virtual nodes to be started as part of the chord ring",
140 Value: 5,
141 Category: "Chord Options",
142 },
143 &cli.StringFlag{
144 Name: "join",
145 Sources: cli.EnvVars("CHORD_JOIN"),
146 Usage: `A known specter server's advertise address.
147 Absent of this flag will bootstrap a new cluster with current node as the seed node`,
148 Category: "Chord Options",
149 },
150 &cli.StringFlag{
151 Name: "kv-provider",
152 Usage: "Backend storage provider for KV. Valid options are memory, aof, and sqlite",
153 Value: "aof",
154 Sources: cli.EnvVars("KV_PROVIDER"),
155 Category: "Chord Options",
156 },
157 
158 &cli.StringFlag{
159 Name: "acme",
160 DefaultText: "acme://{ACME_EMAIL}:@acmehostedzone.com",
161 Sources: cli.EnvVars("ACME_URI"),
162 Usage: `To enable acme, provide an email for the issuer, and the delegated zone for hosting challenges.
163 Absent of this flag will serve self-signed certificate.
164 Alternatively, you can set the URI via the environment variable ACME_URI.`,
165 Category: "Gateway Options",
166 },
167 &cli.IntFlag{
168 Name: "listen-http",
169 Value: 80,
170 Usage: `Override the listening port of the http handler, which handles http connect proxy, and redirects other http requests to https.
171 Note by default the http handler will not be started unless the node is advertising on port 443. Using this option will force the http handler to start.`,
172 Category: "Gateway Options",
173 },
174 &cli.StringSliceFlag{
175 Name: "apex",
176 Sources: cli.EnvVars("APEX"),
177 Usage: "Canonical domain to be used as tunnel root domain. Tunnels will be given names under *.`APEX`. Additional canonical domains can be specified.",
178 Required: true,
179 Category: "Gateway Options",
180 },
181 &cli.StringFlag{
182 Name: "transport-buffer",
183 Value: "16KiB",
184 Usage: "Buffer size when making HTTP request to client",
185 Category: "Gateway Options",
186 },
187 &cli.StringFlag{
188 Name: "proxy-buffer",
189 Value: "16KiB",
190 Usage: "Buffer size when copying response from client",
191 Category: "Gateway Options",
192 },
193 
194 &cli.BoolFlag{
195 Name: "print-acme",
196 Value: false,
197 Usage: "Print acme setup instructions for apex domains based on current configuration.",
198 Category: "Miscellaneous",
199 },
200 
201 &cli.StringFlag{
202 Name: "acme_ca",
203 Hidden: true,
204 Value: cipher.CertCA,
205 Sources: cli.EnvVars("ACME_CA"),
206 },
207 
208 // used for acme setup internally
209 &cli.StringFlag{
210 Name: "acme_email",
211 Hidden: true,
212 },
213 &cli.StringFlag{
214 Name: "acme_zone",
215 Hidden: true,
216 },
217 &cli.StringFlag{
218 Name: "auth_user",
219 Hidden: true,
220 Sources: cli.EnvVars("INTERNAL_USER"),
221 },
222 &cli.StringFlag{
223 Name: "auth_pass",
224 Hidden: true,
225 Sources: cli.EnvVars("INTERNAL_PASS"),
226 },
227 &cli.StringFlag{
228 Name: "env_ca",
229 Hidden: true,
230 Sources: cli.EnvVars("CERT_CA"),
231 },
232 &cli.StringFlag{
233 Name: "env_node",
234 Hidden: true,
235 Sources: cli.EnvVars("CERT_NODE"),
236 },
237 &cli.StringFlag{
238 Name: "env_node_key",
239 Hidden: true,
240 Sources: cli.EnvVars("CERT_NODE_KEY"),
241 },
242 &cli.StringFlag{
243 Name: "env_client_ca",
244 Hidden: true,
245 Sources: cli.EnvVars("CERT_CLIENT_CA"),
246 },
247 &cli.StringFlag{
248 Name: "env_client_ca_key",
249 Hidden: true,
250 Sources: cli.EnvVars("CERT_CLIENT_CA_KEY"),
251 },
252 },
253 Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) {
254 if !cmd.IsSet("cert-dir") && !cmd.IsSet("cert-env") {
255 return ctx, fmt.Errorf("no certificate loader is specified")
256 }
257 if cmd.Int("virtual") < 1 {
258 return ctx, fmt.Errorf("minimum of 1 virtual node is required")
259 }
260 if cmd.IsSet("acme") {
261 email, zone, err := acmeSpec.ParseAcmeURI(cmd.String("acme"))
262 if err != nil {
263 return ctx, err
264 }
265 cmd.Set("acme_email", email)
266 cmd.Set("acme_zone", zone)
267 }
268 return ctx, nil
269 },
270 Action: cmdServer,
271 }
272}
273 
274type certBundle struct {
275 ca *x509.CertPool
276 clientCa *x509.CertPool
277 clientCaCert tls.Certificate
278 node tls.Certificate
279}
280 
281func certLoaderFilesystem(dir string) (*certBundle, error) {
282 files := []string{"ca.crt", "client-ca.crt", "client-ca.key", "node.crt", "node.key"}
283 for i, name := range files {
284 files[i] = filepath.Join(dir, name)
285 }
286 caCert, err := os.ReadFile(files[0])
287 if err != nil {
288 return nil, fmt.Errorf("reading ca cert from file: %w", err)
289 }
290 clientCaCert, err := os.ReadFile(files[1])
291 if err != nil {
292 return nil, fmt.Errorf("reading client ca cert from file: %w", err)
293 }
294 clientCaKey, err := os.ReadFile(files[2])
295 if err != nil {
296 return nil, fmt.Errorf("reading client ca key from file: %w", err)
297 }
298 caCertPool := x509.NewCertPool()
299 if ok := caCertPool.AppendCertsFromPEM(caCert); !ok {
300 return nil, fmt.Errorf("unable to use provided ca bundle")
301 }
302 clientCaCertPool := x509.NewCertPool()
303 if ok := clientCaCertPool.AppendCertsFromPEM(clientCaCert); !ok {
304 return nil, fmt.Errorf("unable to use provided client ca bundle")
305 }
306 node, err := tls.LoadX509KeyPair(files[3], files[4])
307 if err != nil {
308 return nil, fmt.Errorf("creating node cert/key from files: %w", err)
309 }
310 clientCa, err := tls.X509KeyPair(clientCaCert, clientCaKey)
311 if err != nil {
312 return nil, fmt.Errorf("creating client ca cert/key from files: %w", err)
313 }
314 return &certBundle{
315 ca: caCertPool,
316 clientCa: clientCaCertPool,
317 clientCaCert: clientCa,
318 node: node,
319 }, nil
320}
321 
322func certLoaderEnv(cmd *cli.Command) (*certBundle, error) {
323 caCert, err := base64.StdEncoding.DecodeString(cmd.String("env_ca"))
324 if err != nil {
325 return nil, fmt.Errorf("unable to base64 decode CERT_CA: %w", err)
326 }
327 nodeCert, err := base64.StdEncoding.DecodeString(cmd.String("env_node"))
328 if err != nil {
329 return nil, fmt.Errorf("unable to base64 decode CERT_NODE: %w", err)
330 }
331 nodeKey, err := base64.StdEncoding.DecodeString(cmd.String("env_node_key"))
332 if err != nil {
333 return nil, fmt.Errorf("unable to base64 decode CERT_NODE_KEY: %w", err)
334 }
335 clientCaCert, err := base64.StdEncoding.DecodeString(cmd.String("env_client_ca"))
336 if err != nil {
337 return nil, fmt.Errorf("unable to base64 decode CERT_NODE: %w", err)
338 }
339 clientCaKey, err := base64.StdEncoding.DecodeString(cmd.String("env_client_ca_key"))
340 if err != nil {
341 return nil, fmt.Errorf("unable to base64 decode CERT_NODE_KEY: %w", err)
342 }
343 caCertPool := x509.NewCertPool()
344 if ok := caCertPool.AppendCertsFromPEM(caCert); !ok {
345 return nil, fmt.Errorf("unable to use provided ca bundle")
346 }
347 clientCaCertPool := x509.NewCertPool()
348 if ok := clientCaCertPool.AppendCertsFromPEM(clientCaCert); !ok {
349 return nil, fmt.Errorf("unable to use provided client ca bundle")
350 }
351 node, err := tls.X509KeyPair(nodeCert, nodeKey)
352 if err != nil {
353 return nil, fmt.Errorf("creating node cert/key from env: %w", err)
354 }
355 clientCa, err := tls.X509KeyPair(clientCaCert, clientCaKey)
356 if err != nil {
357 return nil, fmt.Errorf("creating client ca cert/key from env: %w", err)
358 }
359 return &certBundle{
360 ca: caCertPool,
361 clientCa: clientCaCertPool,
362 clientCaCert: clientCa,
363 node: node,
364 }, nil
365}
366 
367func configCertProvider(cmd *cli.Command, logger *zap.Logger, kv chord.VNode) (cipher.CertProvider, error) {
368 rootDomains := cmd.StringSlice("apex")
369 managedDomains := rootDomains
370 
371 if cmd.IsSet("acme") {
372 acmeSolver := &acme.ChordSolver{
373 KV: kv,
374 ManagedDomains: managedDomains,
375 }
376 manager, err := acme.NewManager(acme.ManagerConfig{
377 Logger: logger,
378 KV: kv,
379 DNSSolver: acmeSolver,
380 ManagedDomains: managedDomains,
381 CA: cmd.String("acme_ca"),
382 Email: cmd.String("acme_email"),
383 })
384 if err != nil {
385 return nil, err
386 }
387 logger.Info("Using acme as cert provider", zap.String("email", cmd.String("acme_email")), zap.String("zone", cmd.String("acme_zone")))
388 return manager, nil
389 } else {
390 logger.Info("Using self-signed as cert provider")
391 self := &SelfSignedProvider{
392 RootDomain: rootDomains[0],
393 }
394 return self, nil
395 }
396}
397 
398func modifyToSentryLogger(logger *zap.Logger, client *sentry.Client) *zap.Logger {
399 cfg := zapsentry.Configuration{
400 Level: zapcore.WarnLevel,
401 EnableBreadcrumbs: true,
402 BreadcrumbLevel: zapcore.InfoLevel,
403 }
404 core, err := zapsentry.NewCore(cfg, zapsentry.NewSentryClientFromClient(client))
405 
406 if err != nil {
407 logger.Warn("failed to init zap", zap.Error(err))
408 }
409 
410 logger = zapsentry.AttachCoreToLogger(core, logger)
411 
412 return logger
413}
414 
415func cmdServer(ctx context.Context, cmd *cli.Command) error {
416 logger, ok := cmd.Root().Metadata["logger"].(*zap.Logger)
417 if !ok || logger == nil {
418 return fmt.Errorf("unable to obtain logger from app context")
419 }
420 
421 rootDomains := cmd.StringSlice("apex")
422 managedDomains := rootDomains
423 
424 if cmd.Bool("print-acme") {
425 if !cmd.IsSet("acme") {
426 return fmt.Errorf("acme is not configured")
427 }
428 for _, d := range managedDomains {
429 hostname, err := acmeSpec.Normalize(d)
430 if err != nil {
431 return fmt.Errorf("error normalizing domain for %s: %w", d, err)
432 }
433 name, content := acmeSpec.GenerateManagedRecord(hostname, cmd.String("acme_zone"))
434 logger.Info("ACME DNS Record", zap.String("name", name), zap.String("content", content), zap.String("type", "CNAME"))
435 }
436 return nil
437 }
438 
439 dataDir, err := filepath.Abs(cmd.String("data-dir"))
440 if err != nil {
441 return fmt.Errorf("resolving data directory: %w", err)
442 }
443 kvOption := cmd.String("kv-provider")
444 if err := preflightKVProvider(kvOption, dataDir); err != nil {
445 return fmt.Errorf("storage preflight: %w", err)
446 }
447 logger.Info("Storage configuration", zap.String("provider", kvOption), zap.String("data_dir", dataDir))
448 
449 var gwOptions gateway.Options
450 proxyBuffer, err := units.ParseStrictBytes(cmd.String("proxy-buffer"))
451 if err != nil {
452 return fmt.Errorf("error parsing proxy buffer size: %w", err)
453 }
454 gwOptions.ProxyBufferSize = int(proxyBuffer)
455 transportBuffer, err := units.ParseStrictBytes(cmd.String("transport-buffer"))
456 if err != nil {
457 return fmt.Errorf("error parsing transport buffer size: %w", err)
458 }
459 gwOptions.TransportBufferSize = int(transportBuffer)
460 
461 if cmd.IsSet("sentry") {
462 client, err := sentry.NewClient(sentry.ClientOptions{
463 Dsn: cmd.String("sentry"),
464 Release: cmd.Root().Version,
465 })
466 if err != nil {
467 return fmt.Errorf("initializing sentry client: %w", err)
468 }
469 defer client.Flush(time.Second * 2)
470 
471 logger = modifyToSentryLogger(logger, client)
472 defer logger.Sync()
473 }
474 
475 listenBase := cmd.StringSlice("listen-addr")
476 tcpAddrs, err := cmdlisten.ParseAddresses("tcp",
477 listenBase,
478 cmd.StringSlice("listen-tcp"),
479 )
480 if err != nil {
481 return fmt.Errorf("error parsing tcp listen address: %w", err)
482 }
483 
484 udpAddrs, err := cmdlisten.ParseAddresses("udp",
485 listenBase,
486 cmd.StringSlice("listen-udp"),
487 )
488 if err != nil {
489 return fmt.Errorf("error parsing udp listen address: %w", err)
490 }
491 
492 addrStrings := func(addrs []cmdlisten.Address) []string {
493 out := make([]string, 0, len(addrs))
494 for _, a := range addrs {
495 out = append(out, a.Address)
496 }
497 return out
498 }
499 
500 logger.Info("listener configuration",
501 zap.Strings("tcp", addrStrings(tcpAddrs)),
502 zap.Strings("udp", addrStrings(udpAddrs)),
503 )
504 
505 if len(listenBase) == 0 {
506 return fmt.Errorf("at least one listen-addr must be provided")
507 }
508 
509 advertise := listenBase[0]
510 
511 if cmd.IsSet("advertise-addr") {
512 advertise = cmd.String("advertise-addr")
513 }
514 _, advertisePortStr, err := net.SplitHostPort(advertise)
515 if err != nil {
516 return fmt.Errorf("error parsing advertise address: %w", err)
517 }
518 advertisePort, err := strconv.ParseInt(advertisePortStr, 10, 32)
519 if err != nil {
520 return fmt.Errorf("error parsing advertise port: %w", err)
521 }
522 
523 var bundle *certBundle
524 if cmd.IsSet("cert-dir") {
525 bundle, err = certLoaderFilesystem(cmd.String("cert-dir"))
526 if err != nil {
527 return fmt.Errorf("error loading certificates from directory: %w", err)
528 }
529 } else if cmd.IsSet("cert-env") {
530 bundle, err = certLoaderEnv(cmd)
531 if err != nil {
532 return fmt.Errorf("error loading certificates from environment variable: %w", err)
533 }
534 }
535 
536 listenCfg := &net.ListenConfig{
537 Control: reuse.Control,
538 }
539 
540 var (
541 rpcListener net.Listener
542 )
543 if cmd.IsSet("listen-rpc") {
544 parsedRpc, err := url.Parse(cmd.String("listen-rpc"))
545 if err != nil {
546 return fmt.Errorf("error parsing rpc listen address: %w", err)
547 }
548 switch parsedRpc.Scheme {
549 case "unix":
550 rpcListener, err = listenCfg.Listen(ctx, "unix", parsedRpc.Path)
551 case "tcp":
552 rpcListener, err = listenCfg.Listen(ctx, "tcp", parsedRpc.Host)
553 default:
554 return fmt.Errorf("unknown scheme for rpc listen address: %s", parsedRpc.Scheme)
555 }
556 if err != nil {
557 return fmt.Errorf("error setting up rpc listener: %w", err)
558 }
559 defer rpcListener.Close()
560 }
561 
562 tcpListeners := make([]net.Listener, 0, len(tcpAddrs))
563 for _, addr := range tcpAddrs {
564 l, err := listenCfg.Listen(ctx, addr.Network, addr.Address)
565 if err != nil {
566 return fmt.Errorf("error setting up gateway tcp listener on %s: %w", addr.Address, err)
567 }
568 if cmd.Bool("proxy-protocol") {
569 l = &proxyproto.Listener{
570 Listener: l,
571 ReadHeaderTimeout: time.Second * 3,
572 Policy: func(upstream net.Addr) (proxyproto.Policy, error) {
573 return proxyproto.REQUIRE, nil
574 },
575 }
576 }
577 tcpListeners = append(tcpListeners, l)
578 }
579 if len(tcpListeners) == 0 {
580 return fmt.Errorf("no tcp listeners configured")
581 }
582 
583 // TODO: implement SNI proxy so specter can share port with another webserver
584 tcpListener := tcpListeners[0]
585 if len(tcpListeners) > 1 {
586 tcpListener = newMultiListener(tcpListeners)
587 }
588 defer tcpListener.Close()
589 
590 udpBindings := make([]udpBinding, 0, len(udpAddrs))
591 for _, addr := range udpAddrs {
592 pconn, err := listenCfg.ListenPacket(ctx, addr.Network, addr.Address)
593 if err != nil {
594 return fmt.Errorf("error setting up gateway udp listener on %s: %w", addr.Address, err)
595 }
596 tr := &quic.Transport{Conn: pconn}
597 udpBindings = append(udpBindings, udpBinding{
598 listen: addr,
599 packetConn: pconn,
600 transport: tr,
601 })
602 }
603 if len(udpBindings) == 0 {
604 return fmt.Errorf("no udp listeners configured")
605 }
606 for i := range udpBindings {
607 defer udpBindings[i].transport.Close()
608 defer udpBindings[i].packetConn.Close()
609 }
610 
611 var httpListener net.Listener
612 if advertisePort == 443 || cmd.IsSet("listen-http") {
613 httpListeners := make([]net.Listener, 0, len(tcpAddrs))
614 seenHost := make(map[string]struct{}, len(tcpAddrs))
615 for _, addr := range tcpAddrs {
616 if _, ok := seenHost[addr.Host]; ok {
617 continue
618 }
619 seenHost[addr.Host] = struct{}{}
620 httpAddr := net.JoinHostPort(addr.Host, strconv.Itoa(cmd.Int("listen-http")))
621 l, err := listenCfg.Listen(ctx, cmdlisten.NetworkForVersion("tcp", addr.Version), httpAddr)
622 if err != nil {
623 return fmt.Errorf("error setting up http listener on %s: %w", httpAddr, err)
624 }
625 if cmd.Bool("proxy-protocol") {
626 l = &proxyproto.Listener{
627 Listener: l,
628 ReadHeaderTimeout: time.Second * 3,
629 Policy: func(upstream net.Addr) (proxyproto.Policy, error) {
630 return proxyproto.REQUIRE, nil
631 },
632 }
633 }
634 httpListeners = append(httpListeners, l)
635 }
636 
637 if len(httpListeners) > 0 {
638 httpListener = httpListeners[0]
639 if len(httpListeners) > 1 {
640 httpListener = newMultiListener(httpListeners)
641 }
642 defer httpListener.Close()
643 }
644 }
645 
646 alpnMuxes := make([]*overlay.ALPNMux, 0, len(udpBindings))
647 for _, binding := range udpBindings {
648 mux, err := overlay.NewMux(binding.transport)
649 if err != nil {
650 return fmt.Errorf("error setting up quic alpn muxer for %s: %w", binding.listen.Address, err)
651 }
652 alpnMuxes = append(alpnMuxes, mux)
653 defer mux.Close()
654 }
655 
656 chordName := fmt.Sprintf("chord://%s", advertise)
657 tunnelName := fmt.Sprintf("tunnel://%s", advertise)
658 
659 logger.Info("Using advertise addresses as destinations", zap.String("chord", chordName), zap.String("tunnel", tunnelName))
660 
661 chordTLS := cipher.GetPeerTLSConfig(bundle.ca, bundle.node, []string{
662 tun.ALPN(protocol.Link_SPECTER_CHORD),
663 })
664 
665 // handles specter-chord/1
666 chordListeners := make([]q.Listener, 0, len(alpnMuxes))
667 for _, mux := range alpnMuxes {
668 chordListeners = append(chordListeners, mux.With(chordTLS, tun.ALPN(protocol.Link_SPECTER_CHORD)))
669 }
670 chordListener := newMultiQuicListener(ctx, chordListeners)
671 defer chordListener.Close()
672 
673 chordRTT := rtt.NewInstrumentation(20)
674 dialer := newMultiDialer(udpBindings)
675 chordTransport := overlay.NewQUIC(overlay.TransportConfig{
676 Logger: logger.With(zapsentry.NewScope()).With(zap.String("component", "chordTransport")),
677 VirtualTransport: true,
678 ClientTLS: chordTLS,
679 RTTRecorder: chordRTT,
680 QuicTransport: dialer,
681 Endpoint: &protocol.Node{
682 Address: advertise,
683 },
684 })
685 defer chordTransport.Stop()
686 
687 // TODO: measure rtt to client to build routing table with cost
688 tunnelTransport := overlay.NewQUIC(overlay.TransportConfig{
689 UseCertificateIdentity: true,
690 Logger: logger.With(zapsentry.NewScope()).With(zap.String("component", "tunnelTransport")),
691 Endpoint: &protocol.Node{
692 Address: advertise,
693 },
694 QuicTransport: dialer,
695 })
696 defer tunnelTransport.Stop()
697 
698 var existingNode chord.VNode
699 chordClient := rpc.DynamicChordClient(ctx, chordTransport)
700 if cmd.IsSet("join") {
701 existingNode, err = chordImpl.NewRemoteNode(ctx, logger, chordClient, &protocol.Node{
702 Unknown: true,
703 Address: cmd.String("join"),
704 })
705 if err != nil {
706 return fmt.Errorf("error connecting existing chord node: %w", err)
707 }
708 }
709 
710 streamRouter := transport.NewStreamRouter(logger.With(zapsentry.NewScope()).With(zap.String("component", "router")), chordTransport, tunnelTransport)
711 virtualNodes := make([]*chordImpl.LocalNode, 0)
712 
713 k := cmd.Int("virtual")
714 cacheDir := filepath.Join(dataDir, "cache")
715 for i := range k {
716 nodeIdentity := &protocol.Node{
717 Id: chord.Hash(fmt.Appendf(nil, "%s/%d", chordName, i)),
718 Address: advertise,
719 }
720 kvProvider, stopFn, err := getKVProvider(
721 logger.With(zapsentry.NewScope()).With(zap.String("component", "kv"), zap.Object("node", nodeIdentity)),
722 kvOption,
723 filepath.Join(dataDir, fmt.Sprintf("%d", i)),
724 cacheDir,
725 )
726 if err != nil {
727 return fmt.Errorf("initializing kv storage: %w", err)
728 }
729 defer stopFn()
730 
731 virtualNode := chordImpl.NewLocalNode(chordImpl.NodeConfig{
732 Identity: nodeIdentity,
733 BaseLogger: logger,
734 ChordClient: chordClient,
735 KVProvider: kvProvider,
736 StabilizeInterval: timing.ChordStabilizeInterval,
737 FixFingerInterval: timing.ChordFixFingerInterval,
738 PredecessorCheckInterval: timing.ChordPredecessorCheckInterval,
739 NodesRTT: chordRTT,
740 })
741 
742 virtualNode.AttachRouter(ctx, streamRouter)
743 virtualNodes = append(virtualNodes, virtualNode)
744 }
745 
746 rootNode := virtualNodes[0]
747 rootNode.AttachRoot(ctx, streamRouter)
748 if rpcListener != nil {
749 logger.Info("Exposing RPC externally", zap.String("listen", cmd.String("listen-rpc")))
750 rootNode.AttachExternal(ctx, rpcListener)
751 }
752 
753 go chordTransport.AcceptWithListener(ctx, chordListener)
754 go streamRouter.Accept(ctx)
755 for _, mux := range alpnMuxes {
756 go mux.Accept(ctx)
757 }
758 
759 if !cmd.IsSet("join") {
760 if err := rootNode.Create(); err != nil {
761 return fmt.Errorf("error bootstrapping chord ring: %w", err)
762 }
763 } else {
764 if err := rootNode.Join(existingNode); err != nil {
765 return fmt.Errorf("error joining root node to existing chord ring: %w", err)
766 }
767 }
768 defer rootNode.Leave()
769 
770 for i := 1; i < k; i++ {
771 p, err := chordImpl.NewRemoteNode(ctx, logger, chordClient, rootNode.Identity())
772 if err != nil {
773 return fmt.Errorf("error connecting to root node: %w", err)
774 }
775 if err := virtualNodes[i].Join(p); err != nil {
776 return fmt.Errorf("error joining virtual node to root node: %w", err)
777 }
778 defer virtualNodes[i].Leave()
779 }
780 
781 certProvider, err := configCertProvider(cmd, logger.With(zapsentry.NewScope()), chord.WrapRetryKV(rootNode, timing.ChordStabilizeInterval/2, 5))
782 if err != nil {
783 return fmt.Errorf("failed to configure cert provider: %w", err)
784 }
785 
786 if err := certProvider.Initialize(ctx); err != nil {
787 return fmt.Errorf("failed to initialize cert provider: %w", err)
788 }
789 
790 gwTLSConf := cipher.GetGatewayTLSConfig(certProvider.GetCertificate, []string{
791 tun.ALPN(protocol.Link_HTTP2),
792 tun.ALPN(protocol.Link_HTTP),
793 tun.ALPN(protocol.Link_TCP),
794 tun.ALPN(protocol.Link_UNKNOWN),
795 })
796 
797 gwH2Listener := tls.NewListener(tcpListener, gwTLSConf)
798 defer gwH2Listener.Close()
799 
800 // handles h3, h3-29, and specter-tcp/1
801 gwH3Listeners := make([]q.Listener, 0, len(alpnMuxes))
802 for _, mux := range alpnMuxes {
803 gwH3Listeners = append(gwH3Listeners, mux.With(gwTLSConf, append([]string{tun.ALPN(protocol.Link_TCP)}, cipher.H3Protos...)...))
804 }
805 gwH3Listener := newMultiQuicListener(ctx, gwH3Listeners)
806 defer gwH3Listener.Close()
807 
808 // handles specter-client/1
809 clientTLSConf := cipher.GetClientTLSConfig(bundle.clientCa, certProvider.GetCertificate, []string{tun.ALPN(protocol.Link_SPECTER_CLIENT)})
810 clientListeners := make([]q.Listener, 0, len(alpnMuxes))
811 for _, mux := range alpnMuxes {
812 clientListeners = append(clientListeners, mux.With(clientTLSConf, tun.ALPN(protocol.Link_SPECTER_CLIENT)))
813 }
814 clientListener := newMultiQuicListener(ctx, clientListeners)
815 defer clientListener.Close()
816 
817 tunnelIdentity := &protocol.Node{
818 Id: chord.Hash([]byte(tunnelName)),
819 Address: advertise,
820 }
821 tunServer := server.New(server.Config{
822 Logger: logger.With(zapsentry.NewScope()).With(zap.String("component", "tunnelServer"), zap.Uint64("node", tunnelIdentity.GetId())),
823 ParentContext: ctx,
824 Chord: chord.WrapRetryKV(rootNode, timing.ChordStabilizeInterval/2, 5),
825 TunnelTransport: tunnelTransport,
826 ChordTransport: chordTransport,
827 Resolver: net.DefaultResolver,
828 CertProvider: certProvider,
829 Apex: rootDomains[0],
830 Acme: cmd.String("acme_zone"),
831 })
832 defer tunServer.Stop()
833 
834 tunServer.AttachRouter(ctx, streamRouter)
835 tunServer.MustRegister(ctx)
836 
837 go tunnelTransport.AcceptWithListener(ctx, clientListener)
838 
839 var acmeHandler http.Handler
840 if mgr, ok := certProvider.(*acme.Manager); ok {
841 acmeHandler = acme.AcmeManagerHandler(mgr)
842 }
843 gw := gateway.New(gateway.GatewayConfig{
844 PKIServer: &pki.Server{
845 Logger: logger.With(zapsentry.NewScope()).With(zap.String("component", "pki")),
846 ClientCA: bundle.clientCaCert,
847 },
848 Handlers: gateway.InternalHandlers{
849 Acme: acmeHandler,
850 Chord: chordImpl.ChordStatsHandler(rootNode, virtualNodes),
851 Overview: chordImpl.OverviewHandler(rootNode, virtualNodes, kvOption),
852 TunnelServer: server.TunnelServerHandler(tunServer),
853 Migrator: migrator.ConfigMigratorHandler(logger.With(zapsentry.NewScope()).With(zap.String("component", "migrator")), bundle.clientCaCert),
854 },
855 Logger: logger.With(zapsentry.NewScope()).With(zap.String("component", "gateway")),
856 TunnelServer: tunServer,
857 HTTPListener: httpListener,
858 H2Listener: gwH2Listener,
859 H3Listener: gwH3Listener,
860 RootDomains: managedDomains,
861 GatewayPort: int(advertisePort),
862 Options: gwOptions,
863 AdminUser: cmd.String("auth_user"),
864 AdminPass: cmd.String("auth_pass"),
865 HandshakeHintFunc: tunServer.RoutesPreload,
866 })
867 defer gw.Close()
868 
869 gw.AttachRouter(ctx, streamRouter)
870 gw.MustStart(ctx)
871 
872 certProvider.OnHandshake(gw.HandshakeEarlyHint)
873 
874 sigs := make(chan os.Signal, 1)
875 signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM)
876 
877 select {
878 case sig := <-sigs:
879 logger.Info("received signal to stop", zap.String("signal", sig.String()))
880 case <-ctx.Done():
881 logger.Info("context done", zap.Error(ctx.Err()))
882 }
883 
884 return nil
885}