File
Blob: cmd/server/server.go
| 1 | package server |
| 2 | |
| 3 | import ( |
| 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 | |
| 50 | func 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 | |
| 274 | type certBundle struct { |
| 275 | ca *x509.CertPool |
| 276 | clientCa *x509.CertPool |
| 277 | clientCaCert tls.Certificate |
| 278 | node tls.Certificate |
| 279 | } |
| 280 | |
| 281 | func 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 | |
| 322 | func 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 | |
| 367 | func 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 | |
| 398 | func 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 | |
| 415 | func 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 | } |