mirror of
https://github.com/livekit/livekit.git
synced 2026-08-29 07:39:09 +00:00
Remove deprecated egress client (#1701)
* remove deprecated egress client * don't copy mutex
This commit is contained in:
@@ -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"`
|
||||
|
||||
+28
-76
@@ -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
|
||||
}
|
||||
|
||||
+13
-58
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
+1
-11
@@ -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:
|
||||
|
||||
+4
-14
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user