From 7a2eac8e86eaf8c9220b44a1de9dd13753e9dce3 Mon Sep 17 00:00:00 2001 From: Benjamin Pracht Date: Thu, 28 Jul 2022 09:49:54 -0700 Subject: [PATCH] Initial Ingress support in API (#852) This adds support for the Ingress related endpoints to LiveKit server. This currently doesn't handle reconnections safely. --- go.mod | 2 +- go.sum | 4 +- pkg/service/errors.go | 1 + pkg/service/ingress.go | 255 ++++++++++++++++++++++++++++++++++++++ pkg/service/localstore.go | 41 ++---- 5 files changed, 272 insertions(+), 31 deletions(-) create mode 100644 pkg/service/ingress.go diff --git a/go.mod b/go.mod index f0137af2c..73f38661d 100644 --- a/go.mod +++ b/go.mod @@ -15,7 +15,7 @@ require ( github.com/google/wire v0.5.0 github.com/gorilla/websocket v1.4.2 github.com/hashicorp/golang-lru v0.5.4 - github.com/livekit/protocol v0.13.5-0.20220726184153-ad9c55ddef52 + github.com/livekit/protocol v0.13.5-0.20220727215941-ac26418a52e9 github.com/livekit/rtcscore-go v0.0.0-20220524203225-dfd1ba40744a github.com/mackerelio/go-osstat v0.2.1 github.com/magefile/mage v1.13.0 diff --git a/go.sum b/go.sum index 376ffa255..6e3355005 100644 --- a/go.sum +++ b/go.sum @@ -235,10 +235,10 @@ github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/lithammer/shortuuid/v3 v3.0.7 h1:trX0KTHy4Pbwo/6ia8fscyHoGA+mf1jWbPJVuvyJQQ8= github.com/lithammer/shortuuid/v3 v3.0.7/go.mod h1:vMk8ke37EmiewwolSO1NLW8vP4ZaKlRuDIi8tWWmAts= -github.com/livekit/protocol v0.13.5-0.20220721030958-86da2252193b h1:7p5M5WoTFDyIvxcB+r2aqMYKU8tFYf8MA52vXY15naI= -github.com/livekit/protocol v0.13.5-0.20220721030958-86da2252193b/go.mod h1:Qd/Dn4BkJfZQy/IjtEeUOGXARrR7l09WDkg5SY8thkw= github.com/livekit/protocol v0.13.5-0.20220726184153-ad9c55ddef52 h1:E0trQ3RLu2b9hjSiJG1+1hyK/8v57NPJznA7/lKj0qY= github.com/livekit/protocol v0.13.5-0.20220726184153-ad9c55ddef52/go.mod h1:Qd/Dn4BkJfZQy/IjtEeUOGXARrR7l09WDkg5SY8thkw= +github.com/livekit/protocol v0.13.5-0.20220727215941-ac26418a52e9 h1:e12j1EyiiTG56Ag44fwpVtnYQ6MVgLv4bYYI0nTgxZY= +github.com/livekit/protocol v0.13.5-0.20220727215941-ac26418a52e9/go.mod h1:Qd/Dn4BkJfZQy/IjtEeUOGXARrR7l09WDkg5SY8thkw= github.com/livekit/rtcscore-go v0.0.0-20220524203225-dfd1ba40744a h1:cENjhGfslLSDV07gt8ASy47Wd12Q0kBS7hsdunyQ62I= github.com/livekit/rtcscore-go v0.0.0-20220524203225-dfd1ba40744a/go.mod h1:116ych8UaEs9vfIE8n6iZCZ30iagUFTls0vRmC+Ix5U= github.com/mackerelio/go-osstat v0.2.1 h1:5AeAcBEutEErAOlDz6WCkEvm6AKYgHTUQrfwm5RbeQc= diff --git a/pkg/service/errors.go b/pkg/service/errors.go index a464898a0..44d802606 100644 --- a/pkg/service/errors.go +++ b/pkg/service/errors.go @@ -6,6 +6,7 @@ var ( ErrEgressNotFound = errors.New("egress does not exist") ErrEgressNotConnected = errors.New("egress not connected (redis required)") ErrIdentityEmpty = errors.New("identity cannot be empty") + ErrIngressNotConnected = errors.New("ingress not connected (redis required)") ErrIngressNotFound = errors.New("ingress does not exist") ErrMetadataExceedsLimits = errors.New("metadata size exceeds limits") ErrOperationFailed = errors.New("operation cannot be completed") diff --git a/pkg/service/ingress.go b/pkg/service/ingress.go new file mode 100644 index 000000000..fc3db4244 --- /dev/null +++ b/pkg/service/ingress.go @@ -0,0 +1,255 @@ +package service + +import ( + "context" + + "google.golang.org/protobuf/proto" + + "github.com/livekit/livekit-server/pkg/telemetry" + "github.com/livekit/protocol/ingress" + "github.com/livekit/protocol/livekit" + "github.com/livekit/protocol/logger" + "github.com/livekit/protocol/utils" +) + +type IngressService struct { + rpc ingress.RPC + store ServiceStore + roomService livekit.RoomService + telemetry telemetry.TelemetryService + shutdown chan struct{} +} + +func NewIngressService( + rpc ingress.RPC, + store ServiceStore, + rs livekit.RoomService, + ts telemetry.TelemetryService, +) *IngressService { + + return &IngressService{ + rpc: rpc, + store: store, + roomService: rs, + telemetry: ts, + shutdown: make(chan struct{}), + } +} + +func (s *IngressService) Start() { + if s.rpc != nil { + go s.updateWorker() + go s.entitiesWorker() + } +} + +func (s *IngressService) Stop() { + close(s.shutdown) +} + +func (s *IngressService) CreateIngress(ctx context.Context, req *livekit.CreateIngressRequest) (*livekit.IngressInfo, error) { + if err := EnsureRecordPermission(ctx); err != nil { + return nil, twirpAuthError(err) + } + + info := &livekit.IngressInfo{ + IngressId: utils.NewGuid(utils.IngressPrefix), + Name: req.Name, + StreamKey: "TODO", + Url: "TODO", + InputType: req.InputType, + Audio: req.Audio, + Video: req.Video, + RoomName: req.RoomName, + ParticipantIdentity: req.ParticipantIdentity, + ParticipantName: req.ParticipantName, + Reusable: req.InputType == livekit.IngressInput_RTMP_INPUT, + State: &livekit.IngressState{ + Status: livekit.IngressState_ENDPOINT_INACTIVE, + }, + } + + if err := s.store.StoreIngress(ctx, info); err != nil { + logger.Errorw("could not write ingress info", err) + return nil, err + } + + return info, nil +} + +func (s *IngressService) UpdateIngress(ctx context.Context, req *livekit.UpdateIngressRequest) (*livekit.IngressInfo, error) { + if err := EnsureRecordPermission(ctx); err != nil { + return nil, twirpAuthError(err) + } + if s.rpc == nil { + return nil, ErrIngressNotConnected + } + + info, err := s.store.LoadIngress(ctx, req.IngressId) + if err != nil { + logger.Errorw("could not load ingress info", err) + return nil, err + } + + switch info.State.Status { + case livekit.IngressState_ENDPOINT_ERROR: + info.State.Status = livekit.IngressState_ENDPOINT_INACTIVE + fallthrough + + case livekit.IngressState_ENDPOINT_INACTIVE: + if req.Name != "" { + info.Name = req.Name + } + if req.RoomName != "" { + info.RoomName = req.RoomName + } + if req.ParticipantIdentity != "" { + info.ParticipantIdentity = req.ParticipantIdentity + } + if req.ParticipantName != "" { + info.ParticipantName = req.ParticipantName + } + if req.Audio != nil { + info.Audio = req.Audio + } + if req.Video != nil { + info.Video = req.Video + } + + case livekit.IngressState_ENDPOINT_BUFFERING, + livekit.IngressState_ENDPOINT_PUBLISHING: + info, err = s.rpc.SendRequest(ctx, &livekit.IngressRequest{ + IngressId: req.IngressId, + Request: &livekit.IngressRequest_Update{Update: req}, + }) + if err != nil { + logger.Errorw("could not update active ingress", err) + return nil, err + } + } + + err = s.store.UpdateIngress(ctx, info) + if err != nil { + logger.Errorw("could not update ingress info", err) + return nil, err + } + + return info, nil +} + +func (s *IngressService) ListIngress(ctx context.Context, req *livekit.ListIngressRequest) (*livekit.ListIngressResponse, error) { + if err := EnsureRecordPermission(ctx); err != nil { + return nil, twirpAuthError(err) + } + if s.rpc == nil { + return nil, ErrIngressNotConnected + } + + infos, err := s.store.ListIngress(ctx, livekit.RoomName(req.RoomName)) + if err != nil { + logger.Errorw("could not list ingress info", err) + return nil, err + } + + return &livekit.ListIngressResponse{Items: infos}, nil +} + +func (s *IngressService) DeleteIngress(ctx context.Context, req *livekit.DeleteIngressRequest) (*livekit.IngressInfo, error) { + if err := EnsureRecordPermission(ctx); err != nil { + return nil, twirpAuthError(err) + } + if s.rpc == nil { + return nil, ErrIngressNotConnected + } + + info, err := s.store.LoadIngress(ctx, req.IngressId) + if err != nil { + return nil, err + } + + switch info.State.Status { + case livekit.IngressState_ENDPOINT_BUFFERING, + livekit.IngressState_ENDPOINT_PUBLISHING: + info, err = s.rpc.SendRequest(ctx, &livekit.IngressRequest{ + IngressId: req.IngressId, + Request: &livekit.IngressRequest_Delete{Delete: req}, + }) + if err != nil { + logger.Errorw("could not stop active ingress", err) + return nil, err + } + } + + err = s.store.DeleteIngress(ctx, info) + if err != nil { + logger.Errorw("could not delete ingress info", err) + return nil, err + } + + info.State.Status = livekit.IngressState_ENDPOINT_INACTIVE + return info, nil +} + +func (s *IngressService) updateWorker() { + sub, err := s.rpc.GetUpdateChannel(context.Background()) + if err != nil { + logger.Errorw("failed to subscribe to results channel", err) + return + } + + resChan := sub.Channel() + for { + select { + case msg := <-resChan: + b := sub.Payload(msg) + + res := &livekit.IngressInfo{} + if err = proto.Unmarshal(b, res); err != nil { + logger.Errorw("failed to read results", err) + continue + } + + // save updated info to store + err = s.store.UpdateIngress(context.Background(), res) + if err != nil { + logger.Errorw("could not update egress", err) + } + + case <-s.shutdown: + _ = sub.Close() + return + } + } +} + +func (s *IngressService) entitiesWorker() { + sub, err := s.rpc.GetEntityChannel(context.Background()) + if err != nil { + logger.Errorw("failed to subscribe to entities channel", err) + return + } + + resChan := sub.Channel() + for { + select { + case msg := <-resChan: + b := sub.Payload(msg) + + req := &livekit.GetIngressInfoRequest{} + if err = proto.Unmarshal(b, req); err != nil { + logger.Errorw("failed to read request", err) + continue + } + + info, err := s.store.LoadIngress(context.Background(), req.IngressId) + err = s.rpc.SendResponse(context.Background(), req, info, err) + if err != nil { + logger.Errorw("could not send response", err) + } + + case <-s.shutdown: + _ = sub.Close() + return + } + } +} diff --git a/pkg/service/localstore.go b/pkg/service/localstore.go index bb54f0d88..ff5c26986 100644 --- a/pkg/service/localstore.go +++ b/pkg/service/localstore.go @@ -145,63 +145,48 @@ func (s *LocalStore) DeleteParticipant(_ context.Context, roomName livekit.RoomN return nil } +// redis is required for egress func (s *LocalStore) StoreEgress(_ context.Context, _ *livekit.EgressInfo) error { - // redis is required for egress - return nil + return ErrEgressNotConnected } func (s *LocalStore) LoadEgress(_ context.Context, _ string) (*livekit.EgressInfo, error) { - // redis is required for egress - return nil, ErrEgressNotFound + return nil, ErrEgressNotConnected } func (s *LocalStore) ListEgress(_ context.Context, _ livekit.RoomID) ([]*livekit.EgressInfo, error) { - // redis is required for egress - return nil, nil + return nil, ErrEgressNotConnected } func (s *LocalStore) UpdateEgress(_ context.Context, _ *livekit.EgressInfo) error { - // redis is required for egress - return nil + return ErrEgressNotConnected } func (s *LocalStore) DeleteEgress(_ context.Context, _ *livekit.EgressInfo) error { - // redis is required for egress - return nil + return ErrEgressNotConnected } +// redis is required for ingress func (s *LocalStore) StoreIngress(_ context.Context, _ *livekit.IngressInfo) error { - // redis is required for ingress - - return nil + return ErrIngressNotConnected } func (s *LocalStore) LoadIngress(_ context.Context, _ string) (*livekit.IngressInfo, error) { - // redis is required for ingress - - return nil, nil + return nil, ErrIngressNotConnected } func (s *LocalStore) LoadIngressFromStreamKey(_ context.Context, _ string) (*livekit.IngressInfo, error) { - // redis is required for ingress - - return nil, nil + return nil, ErrIngressNotConnected } func (s *LocalStore) ListIngress(_ context.Context, _ livekit.RoomName) ([]*livekit.IngressInfo, error) { - // redis is required for ingress - - return nil, nil + return nil, ErrIngressNotConnected } func (s *LocalStore) UpdateIngress(_ context.Context, _ *livekit.IngressInfo) error { - // redis is required for ingress - - return nil + return ErrIngressNotConnected } func (s *LocalStore) DeleteIngress(_ context.Context, _ *livekit.IngressInfo) error { - // redis is required for ingress - - return nil + return ErrIngressNotConnected }