Initial Ingress support in API (#852)

This adds support for the Ingress related endpoints to LiveKit server. This currently doesn't handle reconnections safely.
This commit is contained in:
Benjamin Pracht
2022-07-28 09:49:54 -07:00
committed by GitHub
parent 188f9c675e
commit 7a2eac8e86
5 changed files with 272 additions and 31 deletions
+1 -1
View File
@@ -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
+2 -2
View File
@@ -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=
+1
View File
@@ -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")
+255
View File
@@ -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
}
}
}
+13 -28
View File
@@ -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
}