diff --git a/README.md b/README.md index 59d9f1cdd..e7098acf2 100644 --- a/README.md +++ b/README.md @@ -91,16 +91,23 @@ APIwLeah7g4fuLYDYAJeaKsSE: 8nTlwISkb-63DPP7OH4e.nw.J44JjicvZDiz8J59EoQ+ In development mode, LiveKit has no external dependencies. With the key file ready, you can start LiveKit with ```shell -./bin/livekit-server --key-file --dev +LIVEKIT_KEYS=": " ./bin/livekit-server --dev ``` or ```shell -docker run --rm -p 7880:7880 -p 7881:7881 -e LIVEKIT_KEYS=": " livekit/livekit-server --dev +docker run --rm \ + -p 7880:7880 \ + -p 7881:7881 \ + -p 7882:7882/udp \ + -e LIVEKIT_KEYS=": " \ + livekit/livekit-server \ + --dev \ + --node-ip= ``` -the `--dev` flag turns on log verbosity to make it easier for local debugging/development +The `--dev` flag turns on log verbosity to make it easier for local debugging/development ### Sample client diff --git a/cmd/server/main.go b/cmd/server/main.go index a32919f06..24674fc3d 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -53,6 +53,11 @@ func main() { Usage: "api keys (key: secret\\n)", EnvVars: []string{"LIVEKIT_KEYS"}, }, + &cli.StringFlag{ + Name: "node-ip", + Usage: "IP address of the current node, used to advertise to clients. Automatically determined by default", + EnvVars: []string{"NODE_IP"}, + }, &cli.StringFlag{ Name: "redis-host", Usage: "host (incl. port) to redis server", @@ -130,15 +135,7 @@ func getConfig(c *cli.Context) (*config.Config, error) { return nil, err } - conf, err := config.NewConfig(confString) - if err != nil { - return nil, err - } - - if err = conf.UpdateFromCLI(c); err != nil { - return nil, err - } - return conf, nil + return config.NewConfig(confString, c) } func startServer(c *cli.Context) error { @@ -177,8 +174,8 @@ func startServer(c *cli.Context) error { defer func() { // run memory profile at termination runtime.GC() - pprof.WriteHeapProfile(f) - f.Close() + _ = pprof.WriteHeapProfile(f) + _ = f.Close() }() } } diff --git a/pkg/config/config.go b/pkg/config/config.go index 4d8352f0b..8bdb86959 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -1,11 +1,13 @@ package config import ( + "fmt" "os" "time" "github.com/mitchellh/go-homedir" "github.com/pion/webrtc/v3" + "github.com/pkg/errors" "github.com/urfave/cli/v2" "gopkg.in/yaml.v3" ) @@ -30,6 +32,7 @@ type RTCConfig struct { TCPPort uint32 `yaml:"tcp_port"` ICEPortRangeStart uint32 `yaml:"port_range_start"` ICEPortRangeEnd uint32 `yaml:"port_range_end"` + NodeIP string `yaml:"node_ip"` // for testing, disable UDP ForceTCP bool `yaml:"force_tcp"` StunServers []string `yaml:"stun_servers"` @@ -90,7 +93,7 @@ type TURNConfig struct { TLSPort int `yaml:"tls_port"` } -func NewConfig(confString string) (*Config, error) { +func NewConfig(confString string, c *cli.Context) (*Config, error) { // start with defaults conf := &Config{ Port: 7880, @@ -98,8 +101,8 @@ func NewConfig(confString string) (*Config, error) { UseExternalIP: false, TCPPort: 7881, UDPPort: 0, - ICEPortRangeStart: 50000, - ICEPortRangeEnd: 60000, + ICEPortRangeStart: 0, + ICEPortRangeEnd: 0, StunServers: []string{ "stun.l.google.com:19302", "stun1.l.google.com:19302", @@ -136,7 +139,39 @@ func NewConfig(confString string) (*Config, error) { Keys: map[string]string{}, } if confString != "" { - _ = yaml.Unmarshal([]byte(confString), conf) + if err := yaml.Unmarshal([]byte(confString), conf); err != nil { + return nil, fmt.Errorf("could not parse config: %v", err) + } + } + + if c != nil { + if err := conf.updateFromCLI(c); err != nil { + return nil, err + } + } + + // expand env vars in filenames + file, err := homedir.Expand(os.ExpandEnv(conf.KeyFile)) + if err != nil { + return nil, err + } + conf.KeyFile = file + + // set defaults for ports if none are set + if conf.RTC.UDPPort == 0 && conf.RTC.ICEPortRangeStart == 0 { + if conf.Development { + conf.RTC.UDPPort = 7882 + } else { + conf.RTC.ICEPortRangeStart = 50000 + conf.RTC.ICEPortRangeEnd = 60000 + } + } + + if conf.RTC.NodeIP == "" { + conf.RTC.NodeIP, err = conf.determineIP() + if err != nil { + return nil, err + } } return conf, nil } @@ -145,7 +180,7 @@ func (conf *Config) HasRedis() bool { return conf.Redis.Address != "" } -func (conf *Config) UpdateFromCLI(c *cli.Context) error { +func (conf *Config) updateFromCLI(c *cli.Context) error { if c.IsSet("dev") { conf.Development = c.Bool("dev") } @@ -154,7 +189,7 @@ func (conf *Config) UpdateFromCLI(c *cli.Context) error { } if c.IsSet("keys") { if err := conf.unmarshalKeys(c.String("keys")); err != nil { - return err + return errors.New("Could not parse keys, it needs to be \"key: secret\", one per line") } } if c.IsSet("redis-host") { @@ -169,12 +204,9 @@ func (conf *Config) UpdateFromCLI(c *cli.Context) error { if c.IsSet("turn-key") { conf.TURN.KeyFile = c.String("turn-key") } - // expand env vars in filenames - file, err := homedir.Expand(os.ExpandEnv(conf.KeyFile)) - if err != nil { - return err + if c.IsSet("node-ip") { + conf.RTC.NodeIP = c.String("node-ip") } - conf.KeyFile = file return nil } diff --git a/pkg/config/config_test.go b/pkg/config/config_test.go index 6f461fe88..35463ef3a 100644 --- a/pkg/config/config_test.go +++ b/pkg/config/config_test.go @@ -7,7 +7,7 @@ import ( ) func TestConfig_UnmarshalKeys(t *testing.T) { - conf, err := NewConfig("") + conf, err := NewConfig("", nil) require.NoError(t, err) require.NoError(t, conf.unmarshalKeys("key1: secret1")) diff --git a/pkg/config/ip.go b/pkg/config/ip.go new file mode 100644 index 000000000..4102a21ac --- /dev/null +++ b/pkg/config/ip.go @@ -0,0 +1,115 @@ +package config + +import ( + "context" + "fmt" + "net" + "time" + + "github.com/livekit/livekit-server/pkg/logger" + "github.com/pion/stun" + "github.com/pkg/errors" +) + +func (conf *Config) determineIP() (string, error) { + if conf.RTC.UseExternalIP { + ip, err := GetExternalIP(conf.RTC.StunServers) + if err == nil { + return ip, nil + } else { + logger.Errorw("could not get external IP", err) + } + } + + // use local ip instead + return GetLocalIPAddress() +} + +func GetLocalIPAddress() (string, error) { + ifaces, err := net.Interfaces() + if err != nil { + return "", err + } + // handle err + var loopBack string + for _, i := range ifaces { + addrs, err := i.Addrs() + if err != nil { + continue + } + for _, addr := range addrs { + var ip net.IP + switch v := addr.(type) { + case *net.IPNet: + ip = v.IP + case *net.IPAddr: + ip = v.IP + default: + continue + } + if ip.IsLoopback() { + loopBack = ip.String() + } else { + return ip.String(), nil + } + } + } + + if loopBack != "" { + return loopBack, nil + } + return "", fmt.Errorf("could not find local IP address") +} + +func GetExternalIP(stunServers []string) (string, error) { + if len(stunServers) == 0 { + return "", errors.New("STUN servers are required but not defined") + } + c, err := stun.Dial("udp4", stunServers[0]) + if err != nil { + return "", err + } + defer c.Close() + + message, err := stun.Build(stun.TransactionID, stun.BindingRequest) + if err != nil { + return "", err + } + + var stunErr error + // sufficiently large buffer to not block it + ipChan := make(chan string, 20) + err = c.Start(message, func(res stun.Event) { + if res.Error != nil { + stunErr = res.Error + return + } + + var xorAddr stun.XORMappedAddress + if err := xorAddr.GetFrom(res.Message); err != nil { + stunErr = err + return + } + ip := xorAddr.IP.To4() + if ip != nil { + ipChan <- ip.String() + } + }) + if err != nil { + return "", err + } + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + select { + case nodeIP := <-ipChan: + return nodeIP, nil + case <-ctx.Done(): + msg := "could not determine public IP" + if stunErr != nil { + return "", errors.Wrap(stunErr, msg) + } else { + return "", fmt.Errorf(msg) + } + } +} diff --git a/pkg/routing/errors.go b/pkg/routing/errors.go index 2a6c82db9..12d6d0056 100644 --- a/pkg/routing/errors.go +++ b/pkg/routing/errors.go @@ -4,6 +4,7 @@ import "errors" var ( ErrNotFound = errors.New("could not find object") + ErrIPNotSet = errors.New("ip address is required and not set") ErrHandlerNotDefined = errors.New("handler not defined") ErrNoAvailableNodes = errors.New("could not find any available nodes") ErrIncorrectRTCNode = errors.New("current node isn't the RTC node for the room") diff --git a/pkg/routing/node.go b/pkg/routing/node.go index 4868e0b60..fff45a901 100644 --- a/pkg/routing/node.go +++ b/pkg/routing/node.go @@ -1,18 +1,13 @@ package routing import ( - "context" "crypto/sha1" "fmt" - "net" "os" "runtime" "time" "github.com/jxskiss/base62" - "github.com/livekit/livekit-server/pkg/logger" - "github.com/pion/stun" - "github.com/pkg/errors" "github.com/livekit/livekit-server/pkg/config" livekit "github.com/livekit/livekit-server/proto" @@ -30,22 +25,16 @@ type NodeStats struct { type LocalNode *livekit.Node func NewLocalNode(conf *config.Config) (LocalNode, error) { - ip, err := GetLocalIP(conf.RTC.StunServers) - if err != nil { - logger.Errorw("could not get local IP", err) - // use local ip instead - ip, err = getLocalIPAddress() - } - if err != nil { - return nil, err - } hostname, err := os.Hostname() if err != nil { return nil, err } + if conf.RTC.NodeIP == "" { + return nil, ErrIPNotSet + } return &livekit.Node{ Id: fmt.Sprintf("%s%s", utils.NodePrefix, HashedID(hostname)[:8]), - Ip: ip, + Ip: conf.RTC.NodeIP, NumCpus: uint32(runtime.NumCPU()), Stats: &livekit.NodeStats{ StartedAt: time.Now().Unix(), @@ -54,59 +43,6 @@ func NewLocalNode(conf *config.Config) (LocalNode, error) { }, nil } -func GetLocalIP(stunServers []string) (string, error) { - if len(stunServers) == 0 { - return "", errors.New("STUN servers are required but not defined") - } - c, err := stun.Dial("udp4", stunServers[0]) - if err != nil { - return "", err - } - defer c.Close() - - message, err := stun.Build(stun.TransactionID, stun.BindingRequest) - if err != nil { - return "", err - } - - var stunErr error - // sufficiently large buffer to not block it - ipChan := make(chan string, 20) - err = c.Start(message, func(res stun.Event) { - if res.Error != nil { - stunErr = res.Error - return - } - - var xorAddr stun.XORMappedAddress - if err := xorAddr.GetFrom(res.Message); err != nil { - stunErr = err - return - } - ip := xorAddr.IP.To4() - if ip != nil { - ipChan <- ip.String() - } - }) - if err != nil { - return "", err - } - - ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) - defer cancel() - select { - case nodeIP := <-ipChan: - return nodeIP, nil - case <-ctx.Done(): - msg := "could not determine public IP" - if stunErr != nil { - return "", errors.Wrap(stunErr, msg) - } else { - return "", fmt.Errorf(msg) - } - } -} - // Creates a hashed ID from a unique string func HashedID(id string) string { h := sha1.New() @@ -115,29 +51,3 @@ func HashedID(id string) string { return base62.EncodeToString(val) } - -func getLocalIPAddress() (string, error) { - ifaces, err := net.Interfaces() - if err != nil { - return "", err - } - // handle err - for _, i := range ifaces { - addrs, err := i.Addrs() - if err != nil { - continue - } - for _, addr := range addrs { - var ip net.IP - switch v := addr.(type) { - case *net.IPNet: - ip = v.IP - case *net.IPAddr: - - ip = v.IP - } - return ip.String(), nil - } - } - return "", fmt.Errorf("could not find local IP address") -} diff --git a/pkg/rtc/config.go b/pkg/rtc/config.go index 41d31eac3..a1ee3df88 100644 --- a/pkg/rtc/config.go +++ b/pkg/rtc/config.go @@ -46,7 +46,7 @@ func NewWebRTCConfig(conf *config.Config, externalIP string) (*WebRTCConfig, err loggerFactory := logging.NewDefaultLoggerFactory() lkLogger := loggerFactory.NewLogger("livekit-mux") - if rtcConf.UseExternalIP && externalIP != "" { + if externalIP != "" { s.SetNAT1To1IPs([]string{externalIP}, webrtc.ICECandidateTypeHost) } diff --git a/pkg/rtc/participant_internal_test.go b/pkg/rtc/participant_internal_test.go index 95307b96d..12ed14a01 100644 --- a/pkg/rtc/participant_internal_test.go +++ b/pkg/rtc/participant_internal_test.go @@ -160,7 +160,7 @@ func TestCorrectJoinedAt(t *testing.T) { } func newParticipantForTest(identity string) *ParticipantImpl { - conf, _ := config.NewConfig("") + conf, _ := config.NewConfig("", nil) // disable mux, it doesn't play too well with unit test conf.RTC.UDPPort = 0 conf.RTC.TCPPort = 0 diff --git a/pkg/service/roommanager_test.go b/pkg/service/roommanager_test.go index b342f126f..03787d6eb 100644 --- a/pkg/service/roommanager_test.go +++ b/pkg/service/roommanager_test.go @@ -27,7 +27,7 @@ func newTestRoomManager(t *testing.T) (*service.RoomManager, *config.Config) { store := &servicefakes.FakeRoomStore{} store.GetRoomReturns(nil, service.ErrRoomNotFound) router := &routingfakes.FakeRouter{} - conf, err := config.NewConfig("") + conf, err := config.NewConfig("", nil) require.NoError(t, err) selector := &routing.RandomSelector{} node, err := routing.NewLocalNode(conf) diff --git a/pkg/service/server.go b/pkg/service/server.go index eaeebc086..a6d56b817 100644 --- a/pkg/service/server.go +++ b/pkg/service/server.go @@ -35,6 +35,7 @@ type LivekitServer struct { currentNode routing.LocalNode running utils.AtomicFlag doneChan chan struct{} + closedChan chan struct{} } func NewLivekitServer(conf *config.Config, @@ -55,6 +56,7 @@ func NewLivekitServer(conf *config.Config, // turn server starts automatically turnServer: turnServer, currentNode: currentNode, + closedChan: make(chan struct{}), } middlewares := []negroni.Handler{ @@ -147,7 +149,8 @@ func (s *LivekitServer) Start() error { go func() { values := []interface{}{ "address", s.httpServer.Addr, - "nodeId", s.currentNode.Id, + "node", s.currentNode.Id, + "nodeIP", s.currentNode.Ip, "version", version.Version, } if s.config.RTC.TCPPort != 0 { @@ -195,6 +198,7 @@ func (s *LivekitServer) Start() error { s.roomManager.Stop() + close(s.closedChan) return nil } @@ -204,8 +208,10 @@ func (s *LivekitServer) Stop() { } s.router.Stop() - s.roomManager.Stop() close(s.doneChan) + + // wait for fully closed + <-s.closedChan } func (s *LivekitServer) RoomManager() *RoomManager { diff --git a/test/integration_helpers.go b/test/integration_helpers.go index 9ff4d8ff1..01f5b865d 100644 --- a/test/integration_helpers.go +++ b/test/integration_helpers.go @@ -49,7 +49,9 @@ func setupSingleNodeTest(name string, roomName string) (*service.LivekitServer, logger.Infow("----------------STARTING TEST----------------", "test", name) s := createSingleNodeServer() go func() { - s.Start() + if err := s.Start(); err != nil { + logger.Errorw("server returned error", err) + } }() waitForServerToStart(s) @@ -131,7 +133,7 @@ func waitUntilConnected(t *testing.T, clients ...*testclient.RTCClient) { func createSingleNodeServer() *service.LivekitServer { var err error - conf, err := config.NewConfig("") + conf, err := config.NewConfig("", nil) if err != nil { panic(fmt.Sprintf("could not create config: %v", err)) } @@ -158,7 +160,7 @@ func createSingleNodeServer() *service.LivekitServer { func createMultiNodeServer(nodeId string, port uint32) *service.LivekitServer { var err error - conf, err := config.NewConfig("") + conf, err := config.NewConfig("", nil) if err != nil { panic(fmt.Sprintf("could not create config: %v", err)) } diff --git a/test/turn_test.go b/test/turn_test.go index 8e66e7369..248766b0f 100644 --- a/test/turn_test.go +++ b/test/turn_test.go @@ -18,7 +18,7 @@ import ( ) func testTurnServer(t *testing.T) { - conf, err := config.NewConfig("") + conf, err := config.NewConfig("", nil) require.NoError(t, err) conf.TURN.Enabled = true