diff --git a/pkg/config/config.go b/pkg/config/config.go index cdab66a70..b379140f6 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -55,7 +55,6 @@ type Config struct { Video VideoConfig `yaml:"video,omitempty"` Room RoomConfig `yaml:"room,omitempty"` TURN TURNConfig `yaml:"turn,omitempty"` - Egress EgressConfig `yaml:"egress,omitempty"` Ingress IngressConfig `yaml:"ingress,omitempty"` WebHook WebHookConfig `yaml:"webhook,omitempty"` NodeSelector NodeSelectorConfig `yaml:"node_selector,omitempty"` @@ -256,10 +255,6 @@ type LimitConfig struct { SubscriptionLimitAudio int32 `yaml:"subscription_limit_audio,omitempty"` } -type EgressConfig struct { - UsePsRPC bool `yaml:"use_psrpc"` -} - type IngressConfig struct { RTMPBaseURL string `yaml:"rtmp_base_url"` WHIPBaseURL string `yaml:"whip_base_url"` diff --git a/pkg/service/egress.go b/pkg/service/egress.go index 80008f82b..99e232b96 100644 --- a/pkg/service/egress.go +++ b/pkg/service/egress.go @@ -10,7 +10,6 @@ import ( "github.com/livekit/livekit-server/pkg/rtc" "github.com/livekit/livekit-server/pkg/telemetry" - "github.com/livekit/protocol/egress" "github.com/livekit/protocol/livekit" "github.com/livekit/protocol/logger" "github.com/livekit/protocol/rpc" @@ -18,42 +17,37 @@ import ( ) type EgressService struct { - psrpcClient rpc.EgressClient - clientDeprecated egress.RPCClient - store ServiceStore - es EgressStore - roomService livekit.RoomService - telemetry telemetry.TelemetryService - launcher rtc.EgressLauncher + client rpc.EgressClient + store ServiceStore + es EgressStore + roomService livekit.RoomService + telemetry telemetry.TelemetryService + launcher rtc.EgressLauncher } type egressLauncher struct { - psrpcClient rpc.EgressClient - clientDeprecated egress.RPCClient - es EgressStore - telemetry telemetry.TelemetryService + client rpc.EgressClient + es EgressStore + telemetry telemetry.TelemetryService } func NewEgressLauncher( - psrpcClient rpc.EgressClient, - clientDeprecated egress.RPCClient, + client rpc.EgressClient, es EgressStore, ts telemetry.TelemetryService) rtc.EgressLauncher { - if psrpcClient == nil && clientDeprecated == nil { + if client == nil { return nil } return &egressLauncher{ - psrpcClient: psrpcClient, - clientDeprecated: clientDeprecated, - es: es, - telemetry: ts, + client: client, + es: es, + telemetry: ts, } } func NewEgressService( - psrpcClient rpc.EgressClient, - clientDeprecated egress.RPCClient, + client rpc.EgressClient, store ServiceStore, es EgressStore, rs livekit.RoomService, @@ -61,13 +55,12 @@ func NewEgressService( launcher rtc.EgressLauncher, ) *EgressService { return &EgressService{ - psrpcClient: psrpcClient, - clientDeprecated: clientDeprecated, - store: store, - es: es, - roomService: rs, - telemetry: ts, - launcher: launcher, + client: client, + store: store, + es: es, + roomService: rs, + telemetry: ts, + launcher: launcher, } } @@ -175,21 +168,12 @@ func (s *egressLauncher) StartEgress(ctx context.Context, req *rpc.StartEgressRe return s.StartEgressWithClusterId(ctx, "", req) } func (s *egressLauncher) StartEgressWithClusterId(ctx context.Context, clusterId string, req *rpc.StartEgressRequest) (*livekit.EgressInfo, error) { - var info *livekit.EgressInfo - var err error - // Ensure we have an Egress ID if req.EgressId == "" { req.EgressId = utils.NewGuid(utils.EgressPrefix) } - if s.psrpcClient != nil { - info, err = s.psrpcClient.StartEgress(ctx, clusterId, req) - } else { - logger.Warnw("Using deprecated egress client. Upgrade egress to v1.5.6+ and use egress:use_psrpc:true in your livekit config", nil) - // SendRequest will transform rpc.StartEgressRequest into deprecated livekit.StartEgressRequest - info, err = s.clientDeprecated.SendRequest(ctx, req) - } + info, err := s.client.StartEgress(ctx, clusterId, req) if err != nil { return nil, err } @@ -213,7 +197,7 @@ func (s *EgressService) UpdateLayout(ctx context.Context, req *livekit.UpdateLay if err := EnsureRecordPermission(ctx); err != nil { return nil, twirpAuthError(err) } - if s.psrpcClient == nil && s.clientDeprecated == nil { + if s.client == nil { return nil, ErrEgressNotConnected } @@ -249,27 +233,11 @@ func (s *EgressService) UpdateStream(ctx context.Context, req *livekit.UpdateStr return nil, twirpAuthError(err) } - if s.psrpcClient == nil && s.clientDeprecated == nil { + if s.client == nil { return nil, ErrEgressNotConnected } - race := rpc.NewRace[livekit.EgressInfo](ctx) - if s.clientDeprecated != nil { - race.Go(func(ctx context.Context) (*livekit.EgressInfo, error) { - return s.clientDeprecated.SendRequest(ctx, &livekit.EgressRequest{ - EgressId: req.EgressId, - Request: &livekit.EgressRequest_UpdateStream{ - UpdateStream: req, - }, - }) - }) - } - if s.psrpcClient != nil { - race.Go(func(ctx context.Context) (*livekit.EgressInfo, error) { - return s.psrpcClient.UpdateStream(ctx, req.EgressId, req) - }) - } - _, info, err := race.Wait() + info, err := s.client.UpdateStream(ctx, req.EgressId, req) if err != nil { return nil, err } @@ -290,7 +258,7 @@ func (s *EgressService) ListEgress(ctx context.Context, req *livekit.ListEgressR if err := EnsureRecordPermission(ctx); err != nil { return nil, twirpAuthError(err) } - if s.psrpcClient == nil && s.clientDeprecated == nil { + if s.client == nil { return nil, ErrEgressNotConnected } @@ -321,7 +289,7 @@ func (s *EgressService) StopEgress(ctx context.Context, req *livekit.StopEgressR return nil, twirpAuthError(err) } - if s.psrpcClient == nil && s.clientDeprecated == nil { + if s.client == nil { return nil, ErrEgressNotConnected } @@ -335,23 +303,7 @@ func (s *EgressService) StopEgress(ctx context.Context, req *livekit.StopEgressR } } - race := rpc.NewRace[livekit.EgressInfo](ctx) - if s.clientDeprecated != nil { - race.Go(func(ctx context.Context) (*livekit.EgressInfo, error) { - return s.clientDeprecated.SendRequest(ctx, &livekit.EgressRequest{ - EgressId: req.EgressId, - Request: &livekit.EgressRequest_Stop{ - Stop: req, - }, - }) - }) - } - if s.psrpcClient != nil { - race.Go(func(ctx context.Context) (*livekit.EgressInfo, error) { - return s.psrpcClient.StopEgress(ctx, req.EgressId, req) - }) - } - _, info, err = race.Wait() + info, err = s.client.StopEgress(ctx, req.EgressId, req) if err != nil { return nil, err } diff --git a/pkg/service/ioinfo.go b/pkg/service/ioinfo.go index e8fc77f9c..529eec8f9 100644 --- a/pkg/service/ioinfo.go +++ b/pkg/service/ioinfo.go @@ -4,11 +4,9 @@ import ( "context" "errors" - "google.golang.org/protobuf/proto" "google.golang.org/protobuf/types/known/emptypb" "github.com/livekit/livekit-server/pkg/telemetry" - "github.com/livekit/protocol/egress" "github.com/livekit/protocol/livekit" "github.com/livekit/protocol/logger" "github.com/livekit/protocol/rpc" @@ -16,12 +14,11 @@ import ( ) type IOInfoService struct { - psrpcServer rpc.IOInfoServer - es EgressStore - is IngressStore - telemetry telemetry.TelemetryService - ecDeprecated egress.RPCClient - shutdown chan struct{} + ioServer rpc.IOInfoServer + es EgressStore + is IngressStore + telemetry telemetry.TelemetryService + shutdown chan struct{} } func NewIOInfoService( @@ -30,22 +27,20 @@ func NewIOInfoService( es EgressStore, is IngressStore, ts telemetry.TelemetryService, - ec egress.RPCClient, ) (*IOInfoService, error) { s := &IOInfoService{ - es: es, - is: is, - telemetry: ts, - ecDeprecated: ec, - shutdown: make(chan struct{}), + es: es, + is: is, + telemetry: ts, + shutdown: make(chan struct{}), } if bus != nil { - psrpcServer, err := rpc.NewIOInfoServer(string(nodeID), s, bus) + ioServer, err := rpc.NewIOInfoServer(string(nodeID), s, bus) if err != nil { return nil, err } - s.psrpcServer = psrpcServer + s.ioServer = ioServer } return s, nil @@ -59,8 +54,6 @@ func (s *IOInfoService) Start() error { logger.Errorw("failed to start redis egress worker", err) return err } - - go s.egressWorkerDeprecated() } return nil @@ -126,45 +119,7 @@ func (s *IOInfoService) UpdateIngressState(ctx context.Context, req *rpc.UpdateI func (s *IOInfoService) Stop() { close(s.shutdown) - if s.psrpcServer != nil { - s.psrpcServer.Shutdown() + if s.ioServer != nil { + s.ioServer.Shutdown() } } - -// Deprecated -func (s *IOInfoService) egressWorkerDeprecated() error { - if s.ecDeprecated == nil { - return nil - } - - go func() { - sub, err := s.ecDeprecated.GetUpdateChannel(context.Background()) - if err != nil { - logger.Errorw("failed to subscribe to results channel", err) - } - - resChan := sub.Channel() - for { - select { - case msg := <-resChan: - b := sub.Payload(msg) - info := &livekit.EgressInfo{} - if err = proto.Unmarshal(b, info); err != nil { - logger.Errorw("failed to read results", err) - continue - } - _, err = s.UpdateEgressInfo(context.Background(), info) - if err != nil { - logger.Errorw("failed to update egress info", err) - } - - case <-s.shutdown: - _ = sub.Close() - s.es.(*RedisStore).Stop() - return - } - } - }() - - return nil -} diff --git a/pkg/service/redisstore.go b/pkg/service/redisstore.go index e7a615ba4..a11c318ad 100644 --- a/pkg/service/redisstore.go +++ b/pkg/service/redisstore.go @@ -505,11 +505,10 @@ func (s *RedisStore) storeIngress(_ context.Context, info *livekit.IngressInfo) } // ignore state - infoCopy := livekit.IngressInfo{} - infoCopy = *info + infoCopy := proto.Clone(info).(*livekit.IngressInfo) infoCopy.State = nil - data, err := proto.Marshal(&infoCopy) + data, err := proto.Marshal(infoCopy) if err != nil { return err } diff --git a/pkg/service/wire.go b/pkg/service/wire.go index 5da0e2c79..126542780 100644 --- a/pkg/service/wire.go +++ b/pkg/service/wire.go @@ -18,7 +18,6 @@ import ( "github.com/livekit/livekit-server/pkg/routing" "github.com/livekit/livekit-server/pkg/telemetry" "github.com/livekit/protocol/auth" - "github.com/livekit/protocol/egress" "github.com/livekit/protocol/livekit" redisLiveKit "github.com/livekit/protocol/redis" "github.com/livekit/protocol/rpc" @@ -45,8 +44,7 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live telemetry.NewTelemetryService, getMessageBus, NewIOInfoService, - getEgressClient, - egress.NewRedisRPCClient, + rpc.NewEgressClient, getEgressStore, NewEgressLauncher, NewEgressService, @@ -148,14 +146,6 @@ func getMessageBus(rc redis.UniversalClient) psrpc.MessageBus { return psrpc.NewRedisMessageBus(rc) } -func getEgressClient(conf *config.Config, nodeID livekit.NodeID, bus psrpc.MessageBus) (rpc.EgressClient, error) { - if conf.Egress.UsePsRPC { - return rpc.NewEgressClient(nodeID, bus) - } - - return nil, nil -} - func getEgressStore(s ObjectStore) EgressStore { switch store := s.(type) { case *RedisStore: diff --git a/pkg/service/wire_gen.go b/pkg/service/wire_gen.go index bce220b63..1ce19893a 100644 --- a/pkg/service/wire_gen.go +++ b/pkg/service/wire_gen.go @@ -13,7 +13,6 @@ import ( "github.com/livekit/livekit-server/pkg/routing" "github.com/livekit/livekit-server/pkg/telemetry" "github.com/livekit/protocol/auth" - "github.com/livekit/protocol/egress" "github.com/livekit/protocol/livekit" redis2 "github.com/livekit/protocol/redis" "github.com/livekit/protocol/rpc" @@ -53,11 +52,10 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live if err != nil { return nil, err } - egressClient, err := getEgressClient(conf, nodeID, messageBus) + egressClient, err := rpc.NewEgressClient(nodeID, messageBus) if err != nil { return nil, err } - rpcClient := egress.NewRedisRPCClient(nodeID, universalClient) egressStore := getEgressStore(objectStore) keyProvider, err := createKeyProvider(conf) if err != nil { @@ -69,12 +67,12 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live } analyticsService := telemetry.NewAnalyticsService(conf, currentNode) telemetryService := telemetry.NewTelemetryService(queuedNotifier, analyticsService) - rtcEgressLauncher := NewEgressLauncher(egressClient, rpcClient, egressStore, telemetryService) + rtcEgressLauncher := NewEgressLauncher(egressClient, egressStore, telemetryService) roomService, err := NewRoomService(roomConfig, apiConfig, router, roomAllocator, objectStore, rtcEgressLauncher) if err != nil { return nil, err } - egressService := NewEgressService(egressClient, rpcClient, objectStore, egressStore, roomService, telemetryService, rtcEgressLauncher) + egressService := NewEgressService(egressClient, objectStore, egressStore, roomService, telemetryService, rtcEgressLauncher) ingressConfig := getIngressConfig(conf) ingressClient, err := rpc.NewIngressClient(nodeID, messageBus) if err != nil { @@ -82,7 +80,7 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live } ingressStore := getIngressStore(objectStore) ingressService := NewIngressService(ingressConfig, nodeID, messageBus, ingressClient, ingressStore, roomService, telemetryService) - ioInfoService, err := NewIOInfoService(nodeID, messageBus, egressStore, ingressStore, telemetryService, rpcClient) + ioInfoService, err := NewIOInfoService(nodeID, messageBus, egressStore, ingressStore, telemetryService) if err != nil { return nil, err } @@ -193,14 +191,6 @@ func getMessageBus(rc redis.UniversalClient) psrpc.MessageBus { return psrpc.NewRedisMessageBus(rc) } -func getEgressClient(conf *config.Config, nodeID livekit.NodeID, bus psrpc.MessageBus) (rpc.EgressClient, error) { - if conf.Egress.UsePsRPC { - return rpc.NewEgressClient(nodeID, bus) - } - - return nil, nil -} - func getEgressStore(s ObjectStore) EgressStore { switch store := s.(type) { case *RedisStore: