IOInfo service (#1305)

* IOInfo service

* only start if not nil

* use ctx in updateEgressInfo

* updates

* fix merge
This commit is contained in:
David Colburn
2023-01-16 16:26:03 -08:00
committed by GitHub
parent f13f0cb52e
commit a87107a0f3
16 changed files with 440 additions and 536 deletions
+3 -3
View File
@@ -19,7 +19,7 @@ require (
github.com/livekit/mageutil v0.0.0-20221221221243-f361fbe40290
github.com/livekit/mediatransportutil v0.0.0-20230111071722-904079e94a7c
github.com/livekit/protocol v1.3.2
github.com/livekit/psrpc v0.2.1
github.com/livekit/psrpc v0.2.3
github.com/livekit/rtcscore-go v0.0.0-20220815072451-20ee10ae1995
github.com/mackerelio/go-osstat v0.2.3
github.com/magefile/mage v1.14.0
@@ -100,8 +100,8 @@ require (
golang.org/x/sys v0.3.0 // indirect
golang.org/x/text v0.5.0 // indirect
golang.org/x/tools v0.1.12 // indirect
google.golang.org/genproto v0.0.0-20200825200019-8632dd797987 // indirect
google.golang.org/grpc v1.51.0 // indirect
google.golang.org/genproto v0.0.0-20221118155620-16455021b5e6 // indirect
google.golang.org/grpc v1.52.0 // indirect
gopkg.in/square/go-jose.v2 v2.6.0 // indirect
gopkg.in/yaml.v2 v2.4.0 // indirect
)
+6 -5
View File
@@ -235,8 +235,8 @@ github.com/livekit/mediatransportutil v0.0.0-20230111071722-904079e94a7c h1:wdzw
github.com/livekit/mediatransportutil v0.0.0-20230111071722-904079e94a7c/go.mod h1:1Dlx20JPoIKGP45eo+yuj0HjeE25zmyeX/EWHiPCjFw=
github.com/livekit/protocol v1.3.2 h1:3goGWbB5HFRb3tMjog8KP0nvZL1Fy6zut3W1psBzqE4=
github.com/livekit/protocol v1.3.2/go.mod h1:gwCG03nKlHlC9hTjL4pXQpn783ALhmbyhq65UZxqbb8=
github.com/livekit/psrpc v0.2.1 h1:ph/4egUMueUPoh5PZ/Aw4v6SH3wAbA+2t/GyCbpPKTg=
github.com/livekit/psrpc v0.2.1/go.mod h1:MCe0xLdFPXmzogPiLrM94JIJbctb9+fAv5qYPkY2DXw=
github.com/livekit/psrpc v0.2.3 h1:M8B/BdrrpRUj3Uj2MeOIDvJDFqiydoN2yeW4P0EF4ZQ=
github.com/livekit/psrpc v0.2.3/go.mod h1:+nJvbKx9DCZ6PSAsMHJPRAKjmRJ5WiyyhEmbKYqMKto=
github.com/livekit/rtcscore-go v0.0.0-20220815072451-20ee10ae1995 h1:vOaY2qvfLihDyeZtnGGN1Law9wRrw8BMGCr1TygTvMw=
github.com/livekit/rtcscore-go v0.0.0-20220815072451-20ee10ae1995/go.mod h1:116ych8UaEs9vfIE8n6iZCZ30iagUFTls0vRmC+Ix5U=
github.com/mackerelio/go-osstat v0.2.3 h1:jAMXD5erlDE39kdX2CU7YwCGRcxIO33u/p8+Fhe5dJw=
@@ -735,8 +735,9 @@ google.golang.org/genproto v0.0.0-20200526211855-cb27e3aa2013/go.mod h1:NbSheEEY
google.golang.org/genproto v0.0.0-20200618031413-b414f8b61790/go.mod h1:jDfRM7FcilCzHH/e9qn6dsT145K34l5v+OpcnNgKAAA=
google.golang.org/genproto v0.0.0-20200729003335-053ba62fc06f/go.mod h1:FWY/as6DDZQgahTzZj3fqbO1CbirC29ZNUFHwi0/+no=
google.golang.org/genproto v0.0.0-20200804131852-c06518451d9c/go.mod h1:FWY/as6DDZQgahTzZj3fqbO1CbirC29ZNUFHwi0/+no=
google.golang.org/genproto v0.0.0-20200825200019-8632dd797987 h1:PDIOdWxZ8eRizhKa1AAvY53xsvLB1cWorMjslvY3VA8=
google.golang.org/genproto v0.0.0-20200825200019-8632dd797987/go.mod h1:FWY/as6DDZQgahTzZj3fqbO1CbirC29ZNUFHwi0/+no=
google.golang.org/genproto v0.0.0-20221118155620-16455021b5e6 h1:a2S6M0+660BgMNl++4JPlcAO/CjkqYItDEZwkoDQK7c=
google.golang.org/genproto v0.0.0-20221118155620-16455021b5e6/go.mod h1:rZS5c/ZVYMaOGBfO68GWtjOw/eLaZM1X6iVtgjZ+EWg=
google.golang.org/grpc v1.19.0/go.mod h1:mqu4LbDTu4XGKhr4mRzUsmM4RtVoemTSY81AxZiDr8c=
google.golang.org/grpc v1.20.1/go.mod h1:10oTOabMzJvdu6/UiuZezV6QK5dSlG84ov/aaiqXj38=
google.golang.org/grpc v1.21.1/go.mod h1:oYelfM1adQP15Ek0mdvEgi9Df8B9CZIaU1084ijfRaM=
@@ -749,8 +750,8 @@ google.golang.org/grpc v1.28.0/go.mod h1:rpkK4SK4GF4Ach/+MFLZUBavHOvF2JJB5uozKKa
google.golang.org/grpc v1.29.1/go.mod h1:itym6AZVZYACWQqET3MqgPpjcuV5QH3BxFS3IjizoKk=
google.golang.org/grpc v1.30.0/go.mod h1:N36X2cJ7JwdamYAgDz+s+rVMFjt3numwzf/HckM8pak=
google.golang.org/grpc v1.31.0/go.mod h1:N36X2cJ7JwdamYAgDz+s+rVMFjt3numwzf/HckM8pak=
google.golang.org/grpc v1.51.0 h1:E1eGv1FTqoLIdnBCZufiSHgKjlqG6fKFf6pPWtMTh8U=
google.golang.org/grpc v1.51.0/go.mod h1:wgNDFcnuBGmxLKI/qn4T+m5BtEBYXJPvibbUPsAIPww=
google.golang.org/grpc v1.52.0 h1:kd48UiU7EHsV4rnLyOJRuP/Il/UHE7gdDAQ+SZI7nZk=
google.golang.org/grpc v1.52.0/go.mod h1:pu6fVzoFb+NBYNAvQL08ic+lvB2IojljRYuun5vorUY=
google.golang.org/protobuf v0.0.0-20200109180630-ec00e32a8dfd/go.mod h1:DFci5gLYBciE7Vtevhsrf46CRTquxDuWsQurQQe4oz8=
google.golang.org/protobuf v0.0.0-20200221191635-4d8936d0db64/go.mod h1:kwYJMbMJ01Woi6D6+Kah6886xMZcty6N08ah7+eCXa0=
google.golang.org/protobuf v0.0.0-20200228230310-ab0ca4ff8a60/go.mod h1:cfTl7dwQJ+fmap5saPgwCLgHXTUD7jkjRqWcaiX5VyM=
-110
View File
@@ -3,12 +3,8 @@ package service
import (
"context"
"encoding/json"
"errors"
"fmt"
"reflect"
"time"
"google.golang.org/protobuf/proto"
"github.com/livekit/livekit-server/pkg/rtc"
"github.com/livekit/livekit-server/pkg/service/rpc"
@@ -28,7 +24,6 @@ type EgressService struct {
roomService livekit.RoomService
telemetry telemetry.TelemetryService
launcher rtc.EgressLauncher
shutdown chan struct{}
}
type egressLauncher struct {
@@ -75,23 +70,6 @@ func NewEgressService(
}
}
func (s *EgressService) Start() error {
if s.shutdown != nil {
return nil
}
s.shutdown = make(chan struct{})
if (s.psrpcClient != nil || s.clientDeprecated != nil) && s.es != nil {
return s.startWorker()
}
return nil
}
func (s *EgressService) Stop() {
close(s.shutdown)
}
func (s *EgressService) StartRoomCompositeEgress(ctx context.Context, req *livekit.RoomCompositeEgressRequest) (*livekit.EgressInfo, error) {
fields := []interface{}{"room", req.RoomName, "outputType", reflect.TypeOf(req.Output).String(), "baseUrl", req.CustomBaseUrl}
defer func() {
@@ -356,91 +334,3 @@ func (s *EgressService) StopEgress(ctx context.Context, req *livekit.StopEgressR
return info, nil
}
func (s *EgressService) startWorker() error {
rs := s.es.(*RedisStore)
err := rs.Start()
if err != nil {
logger.Errorw("failed to start redis egress worker", err)
return err
}
if s.psrpcClient != nil {
go func() {
sub, err := s.psrpcClient.SubscribeInfoUpdate(context.Background())
if err != nil {
logger.Errorw("failed to subscribe", err)
}
for {
select {
case info := <-sub.Channel():
s.handleUpdate(info)
case <-s.shutdown:
_ = sub.Close()
return
}
}
}()
}
if s.clientDeprecated != nil {
go func() {
sub, err := s.clientDeprecated.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
}
s.handleUpdate(info)
case <-s.shutdown:
_ = sub.Close()
rs.Stop()
return
}
}
}()
}
return nil
}
func (s *EgressService) handleUpdate(info *livekit.EgressInfo) {
switch info.Status {
case livekit.EgressStatus_EGRESS_COMPLETE,
livekit.EgressStatus_EGRESS_FAILED,
livekit.EgressStatus_EGRESS_ABORTED:
// make sure endedAt is set so it eventually gets deleted
if info.EndedAt == 0 {
info.EndedAt = time.Now().UnixNano()
}
if err := s.es.UpdateEgress(context.Background(), info); err != nil {
logger.Errorw("could not update egress", err)
}
// log results
if info.Error != "" {
logger.Errorw("egress failed", errors.New(info.Error), "egressID", info.EgressId)
} else {
logger.Infow("egress ended", "egressID", info.EgressId)
}
s.telemetry.EgressEnded(context.Background(), info)
default:
if err := s.es.UpdateEgress(context.Background(), info); err != nil {
logger.Errorw("could not update egress", err)
}
}
}
-123
View File
@@ -2,11 +2,8 @@ package service
import (
"context"
"errors"
"time"
"google.golang.org/protobuf/proto"
"github.com/livekit/livekit-server/pkg/config"
"github.com/livekit/livekit-server/pkg/service/rpc"
"github.com/livekit/livekit-server/pkg/telemetry"
@@ -27,12 +24,10 @@ type IngressService struct {
nodeID livekit.NodeID
bus psrpc.MessageBus
psrpcClient rpc.IngressClient
psrpcServer rpc.IOInfoServer
rpcClient ingress.RPCClient
store IngressStore
roomService livekit.RoomService
telemetry telemetry.TelemetryService
shutdown chan struct{}
}
func NewIngressService(
@@ -55,29 +50,6 @@ func NewIngressService(
store: store,
roomService: rs,
telemetry: ts,
shutdown: make(chan struct{}),
}
}
func (s *IngressService) Start() error {
if s.psrpcClient != nil {
psrpcServer, err := rpc.NewIOInfoServer(string(s.nodeID), s, s.bus)
if err != nil {
return err
}
s.psrpcServer = psrpcServer
} else if s.rpcClient != nil {
go s.updateWorker()
go s.entitiesWorker()
}
return nil
}
func (s *IngressService) Stop() {
close(s.shutdown)
if s.psrpcServer != nil {
s.psrpcServer.Shutdown()
}
}
@@ -320,98 +292,3 @@ func (s *IngressService) DeleteIngress(ctx context.Context, req *livekit.DeleteI
info.State.Status = livekit.IngressState_ENDPOINT_INACTIVE
return info, nil
}
func (s *IngressService) updateWorker() {
sub, err := s.rpcClient.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.UpdateIngressStateRequest{}
if err = proto.Unmarshal(b, res); err != nil {
logger.Errorw("failed to read results", err)
continue
}
// save updated info to store
err = s.store.UpdateIngressState(context.Background(), res.IngressId, res.State)
if err != nil {
logger.Errorw("could not update ingress", err)
}
case <-s.shutdown:
_ = sub.Close()
return
}
}
}
func (s *IngressService) UpdateIngressState(ctx context.Context, req *livekit.UpdateIngressStateRequest) (*rpc.Ignored, error) {
if err := s.store.UpdateIngressState(ctx, req.IngressId, req.State); err != nil {
logger.Errorw("could not update ingress", err)
return nil, err
}
return &rpc.Ignored{}, nil
}
func (s *IngressService) loadIngressFromInfoRequest(req *livekit.GetIngressInfoRequest) (info *livekit.IngressInfo, err error) {
if req.IngressId != "" {
info, err = s.store.LoadIngress(context.Background(), req.IngressId)
} else if req.StreamKey != "" {
info, err = s.store.LoadIngressFromStreamKey(context.Background(), req.StreamKey)
} else {
err = errors.New("request needs to specity either IngressId or StreamKey")
}
return info, err
}
func (s *IngressService) entitiesWorker() {
sub, err := s.rpcClient.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.loadIngressFromInfoRequest(req)
if err != nil {
logger.Errorw("failed to load ingress info", err)
continue
}
err = s.rpcClient.SendGetIngressInfoResponse(context.Background(), req, &livekit.GetIngressInfoResponse{Info: info}, err)
if err != nil {
logger.Errorw("could not send response", err)
}
case <-s.shutdown:
_ = sub.Close()
return
}
}
}
func (s *IngressService) GetIngressInfo(ctx context.Context, req *livekit.GetIngressInfoRequest) (*livekit.GetIngressInfoResponse, error) {
info, err := s.loadIngressFromInfoRequest(req)
if err != nil {
return nil, err
}
return &livekit.GetIngressInfoResponse{Info: info}, nil
}
+247
View File
@@ -0,0 +1,247 @@
package service
import (
"context"
"errors"
"time"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/emptypb"
"github.com/livekit/livekit-server/pkg/service/rpc"
"github.com/livekit/livekit-server/pkg/telemetry"
"github.com/livekit/protocol/egress"
"github.com/livekit/protocol/ingress"
"github.com/livekit/protocol/livekit"
"github.com/livekit/protocol/logger"
"github.com/livekit/psrpc"
)
type IOInfoService struct {
psrpcServer rpc.IOInfoServer
es EgressStore
is IngressStore
telemetry telemetry.TelemetryService
ecDeprecated egress.RPCClient
icDeprecated ingress.RPCClient
shutdown chan struct{}
}
func NewIOInfoService(
nodeID livekit.NodeID,
bus psrpc.MessageBus,
es EgressStore,
is IngressStore,
ts telemetry.TelemetryService,
ec egress.RPCClient,
ic ingress.RPCClient,
) (*IOInfoService, error) {
s := &IOInfoService{
es: es,
is: is,
telemetry: ts,
ecDeprecated: ec,
icDeprecated: ic,
shutdown: make(chan struct{}),
}
if bus != nil {
psrpcServer, err := rpc.NewIOInfoServer(string(nodeID), s, bus)
if err != nil {
return nil, err
}
s.psrpcServer = psrpcServer
}
return s, nil
}
func (s *IOInfoService) Start() error {
if s.es != nil {
rs := s.es.(*RedisStore)
err := rs.Start()
if err != nil {
logger.Errorw("failed to start redis egress worker", err)
return err
}
go s.egressWorkerDeprecated()
go s.ingressWorkerDeprecated()
}
return nil
}
func (s *IOInfoService) UpdateEgressInfo(ctx context.Context, info *livekit.EgressInfo) (*emptypb.Empty, error) {
switch info.Status {
case livekit.EgressStatus_EGRESS_COMPLETE,
livekit.EgressStatus_EGRESS_FAILED,
livekit.EgressStatus_EGRESS_ABORTED:
// make sure endedAt is set so it eventually gets deleted
if info.EndedAt == 0 {
info.EndedAt = time.Now().UnixNano()
}
if err := s.es.UpdateEgress(ctx, info); err != nil {
logger.Errorw("could not update egress", err)
return nil, err
}
// log results
if info.Error != "" {
logger.Errorw("egress failed", errors.New(info.Error), "egressID", info.EgressId)
} else {
logger.Infow("egress ended", "egressID", info.EgressId)
}
s.telemetry.EgressEnded(ctx, info)
default:
if err := s.es.UpdateEgress(ctx, info); err != nil {
logger.Errorw("could not update egress", err)
return nil, err
}
}
return &emptypb.Empty{}, nil
}
func (s *IOInfoService) GetIngressInfo(ctx context.Context, req *livekit.GetIngressInfoRequest) (*livekit.GetIngressInfoResponse, error) {
info, err := s.loadIngressFromInfoRequest(req)
if err != nil {
return nil, err
}
return &livekit.GetIngressInfoResponse{Info: info}, nil
}
func (s *IOInfoService) loadIngressFromInfoRequest(req *livekit.GetIngressInfoRequest) (info *livekit.IngressInfo, err error) {
if req.IngressId != "" {
info, err = s.is.LoadIngress(context.Background(), req.IngressId)
} else if req.StreamKey != "" {
info, err = s.is.LoadIngressFromStreamKey(context.Background(), req.StreamKey)
} else {
err = errors.New("request needs to specity either IngressId or StreamKey")
}
return info, err
}
func (s *IOInfoService) UpdateIngressState(ctx context.Context, req *livekit.UpdateIngressStateRequest) (*emptypb.Empty, error) {
if err := s.is.UpdateIngressState(ctx, req.IngressId, req.State); err != nil {
logger.Errorw("could not update ingress", err)
return nil, err
}
return &emptypb.Empty{}, nil
}
func (s *IOInfoService) Stop() {
close(s.shutdown)
if s.psrpcServer != nil {
s.psrpcServer.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
}
// Deprecated
func (s *IOInfoService) ingressWorkerDeprecated() {
if s.icDeprecated == nil {
return
}
updates, err := s.icDeprecated.GetUpdateChannel(context.Background())
if err != nil {
logger.Errorw("failed to subscribe to results channel", err)
return
}
entities, err := s.icDeprecated.GetEntityChannel(context.Background())
if err != nil {
logger.Errorw("failed to subscribe to entities channel", err)
_ = updates.Close()
return
}
updateChan := updates.Channel()
entityChan := entities.Channel()
for {
select {
case msg := <-updateChan:
b := updates.Payload(msg)
res := &livekit.UpdateIngressStateRequest{}
if err = proto.Unmarshal(b, res); err != nil {
logger.Errorw("failed to read results", err)
continue
}
// save updated info to store
err = s.is.UpdateIngressState(context.Background(), res.IngressId, res.State)
if err != nil {
logger.Errorw("could not update ingress", err)
}
case msg := <-entityChan:
b := entities.Payload(msg)
req := &livekit.GetIngressInfoRequest{}
if err = proto.Unmarshal(b, req); err != nil {
logger.Errorw("failed to read request", err)
continue
}
info, err := s.loadIngressFromInfoRequest(req)
if err != nil {
logger.Errorw("failed to load ingress info", err)
continue
}
err = s.icDeprecated.SendGetIngressInfoResponse(context.Background(), req, &livekit.GetIngressInfoResponse{Info: info}, err)
if err != nil {
logger.Errorw("could not send response", err)
}
case <-s.shutdown:
_ = updates.Close()
_ = entities.Close()
return
}
}
}
+53 -110
View File
@@ -1,7 +1,7 @@
// Code generated by protoc-gen-go. DO NOT EDIT.
// versions:
// protoc-gen-go v1.28.1
// protoc v3.21.12
// protoc-gen-go v1.27.1
// protoc v3.21.5
// source: pkg/service/rpc/egress.proto
package rpc
@@ -22,44 +22,6 @@ const (
_ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20)
)
type Empty struct {
state protoimpl.MessageState
sizeCache protoimpl.SizeCache
unknownFields protoimpl.UnknownFields
}
func (x *Empty) Reset() {
*x = Empty{}
if protoimpl.UnsafeEnabled {
mi := &file_pkg_service_rpc_egress_proto_msgTypes[0]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
}
func (x *Empty) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*Empty) ProtoMessage() {}
func (x *Empty) ProtoReflect() protoreflect.Message {
mi := &file_pkg_service_rpc_egress_proto_msgTypes[0]
if protoimpl.UnsafeEnabled && x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use Empty.ProtoReflect.Descriptor instead.
func (*Empty) Descriptor() ([]byte, []int) {
return file_pkg_service_rpc_egress_proto_rawDescGZIP(), []int{0}
}
type ListActiveEgressRequest struct {
state protoimpl.MessageState
sizeCache protoimpl.SizeCache
@@ -69,7 +31,7 @@ type ListActiveEgressRequest struct {
func (x *ListActiveEgressRequest) Reset() {
*x = ListActiveEgressRequest{}
if protoimpl.UnsafeEnabled {
mi := &file_pkg_service_rpc_egress_proto_msgTypes[1]
mi := &file_pkg_service_rpc_egress_proto_msgTypes[0]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -82,7 +44,7 @@ func (x *ListActiveEgressRequest) String() string {
func (*ListActiveEgressRequest) ProtoMessage() {}
func (x *ListActiveEgressRequest) ProtoReflect() protoreflect.Message {
mi := &file_pkg_service_rpc_egress_proto_msgTypes[1]
mi := &file_pkg_service_rpc_egress_proto_msgTypes[0]
if protoimpl.UnsafeEnabled && x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -95,7 +57,7 @@ func (x *ListActiveEgressRequest) ProtoReflect() protoreflect.Message {
// Deprecated: Use ListActiveEgressRequest.ProtoReflect.Descriptor instead.
func (*ListActiveEgressRequest) Descriptor() ([]byte, []int) {
return file_pkg_service_rpc_egress_proto_rawDescGZIP(), []int{1}
return file_pkg_service_rpc_egress_proto_rawDescGZIP(), []int{0}
}
type ListActiveEgressResponse struct {
@@ -109,7 +71,7 @@ type ListActiveEgressResponse struct {
func (x *ListActiveEgressResponse) Reset() {
*x = ListActiveEgressResponse{}
if protoimpl.UnsafeEnabled {
mi := &file_pkg_service_rpc_egress_proto_msgTypes[2]
mi := &file_pkg_service_rpc_egress_proto_msgTypes[1]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
@@ -122,7 +84,7 @@ func (x *ListActiveEgressResponse) String() string {
func (*ListActiveEgressResponse) ProtoMessage() {}
func (x *ListActiveEgressResponse) ProtoReflect() protoreflect.Message {
mi := &file_pkg_service_rpc_egress_proto_msgTypes[2]
mi := &file_pkg_service_rpc_egress_proto_msgTypes[1]
if protoimpl.UnsafeEnabled && x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
@@ -135,7 +97,7 @@ func (x *ListActiveEgressResponse) ProtoReflect() protoreflect.Message {
// Deprecated: Use ListActiveEgressResponse.ProtoReflect.Descriptor instead.
func (*ListActiveEgressResponse) Descriptor() ([]byte, []int) {
return file_pkg_service_rpc_egress_proto_rawDescGZIP(), []int{2}
return file_pkg_service_rpc_egress_proto_rawDescGZIP(), []int{1}
}
func (x *ListActiveEgressResponse) GetEgressIds() []string {
@@ -154,38 +116,34 @@ var file_pkg_service_rpc_egress_proto_rawDesc = []byte{
0x74, 0x6f, 0x1a, 0x1a, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x5f, 0x72, 0x70, 0x63, 0x5f,
0x69, 0x6e, 0x74, 0x65, 0x72, 0x6e, 0x61, 0x6c, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x1a, 0x14,
0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x5f, 0x65, 0x67, 0x72, 0x65, 0x73, 0x73, 0x2e, 0x70,
0x72, 0x6f, 0x74, 0x6f, 0x22, 0x07, 0x0a, 0x05, 0x45, 0x6d, 0x70, 0x74, 0x79, 0x22, 0x19, 0x0a,
0x17, 0x4c, 0x69, 0x73, 0x74, 0x41, 0x63, 0x74, 0x69, 0x76, 0x65, 0x45, 0x67, 0x72, 0x65, 0x73,
0x73, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x22, 0x39, 0x0a, 0x18, 0x4c, 0x69, 0x73, 0x74,
0x41, 0x63, 0x74, 0x69, 0x76, 0x65, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x52, 0x65, 0x73, 0x70,
0x6f, 0x6e, 0x73, 0x65, 0x12, 0x1d, 0x0a, 0x0a, 0x65, 0x67, 0x72, 0x65, 0x73, 0x73, 0x5f, 0x69,
0x64, 0x73, 0x18, 0x01, 0x20, 0x03, 0x28, 0x09, 0x52, 0x09, 0x65, 0x67, 0x72, 0x65, 0x73, 0x73,
0x49, 0x64, 0x73, 0x32, 0xb2, 0x01, 0x0a, 0x0e, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x6e,
0x74, 0x65, 0x72, 0x6e, 0x61, 0x6c, 0x12, 0x47, 0x0a, 0x0b, 0x53, 0x74, 0x61, 0x72, 0x74, 0x45,
0x67, 0x72, 0x65, 0x73, 0x73, 0x12, 0x1b, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e,
0x53, 0x74, 0x61, 0x72, 0x74, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x52, 0x65, 0x71, 0x75, 0x65,
0x73, 0x74, 0x1a, 0x13, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x45, 0x67, 0x72,
0x65, 0x73, 0x73, 0x49, 0x6e, 0x66, 0x6f, 0x22, 0x06, 0xb2, 0x89, 0x01, 0x02, 0x20, 0x01, 0x12,
0x57, 0x0a, 0x10, 0x4c, 0x69, 0x73, 0x74, 0x41, 0x63, 0x74, 0x69, 0x76, 0x65, 0x45, 0x67, 0x72,
0x65, 0x73, 0x73, 0x12, 0x1c, 0x2e, 0x72, 0x70, 0x63, 0x2e, 0x4c, 0x69, 0x73, 0x74, 0x41, 0x63,
0x74, 0x69, 0x76, 0x65, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73,
0x74, 0x1a, 0x1d, 0x2e, 0x72, 0x70, 0x63, 0x2e, 0x4c, 0x69, 0x73, 0x74, 0x41, 0x63, 0x74, 0x69,
0x76, 0x65, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65,
0x22, 0x06, 0xb2, 0x89, 0x01, 0x02, 0x08, 0x01, 0x32, 0xd8, 0x01, 0x0a, 0x0d, 0x45, 0x67, 0x72,
0x65, 0x73, 0x73, 0x48, 0x61, 0x6e, 0x64, 0x6c, 0x65, 0x72, 0x12, 0x49, 0x0a, 0x0c, 0x55, 0x70,
0x64, 0x61, 0x74, 0x65, 0x53, 0x74, 0x72, 0x65, 0x61, 0x6d, 0x12, 0x1c, 0x2e, 0x6c, 0x69, 0x76,
0x65, 0x6b, 0x69, 0x74, 0x2e, 0x55, 0x70, 0x64, 0x61, 0x74, 0x65, 0x53, 0x74, 0x72, 0x65, 0x61,
0x6d, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x13, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b,
0x69, 0x74, 0x2e, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x6e, 0x66, 0x6f, 0x22, 0x06, 0xb2,
0x89, 0x01, 0x02, 0x18, 0x01, 0x12, 0x45, 0x0a, 0x0a, 0x53, 0x74, 0x6f, 0x70, 0x45, 0x67, 0x72,
0x65, 0x73, 0x73, 0x12, 0x1a, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x53, 0x74,
0x6f, 0x70, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a,
0x72, 0x6f, 0x74, 0x6f, 0x22, 0x19, 0x0a, 0x17, 0x4c, 0x69, 0x73, 0x74, 0x41, 0x63, 0x74, 0x69,
0x76, 0x65, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x22,
0x39, 0x0a, 0x18, 0x4c, 0x69, 0x73, 0x74, 0x41, 0x63, 0x74, 0x69, 0x76, 0x65, 0x45, 0x67, 0x72,
0x65, 0x73, 0x73, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x1d, 0x0a, 0x0a, 0x65,
0x67, 0x72, 0x65, 0x73, 0x73, 0x5f, 0x69, 0x64, 0x73, 0x18, 0x01, 0x20, 0x03, 0x28, 0x09, 0x52,
0x09, 0x65, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x64, 0x73, 0x32, 0xb2, 0x01, 0x0a, 0x0e, 0x45,
0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x6e, 0x74, 0x65, 0x72, 0x6e, 0x61, 0x6c, 0x12, 0x47, 0x0a,
0x0b, 0x53, 0x74, 0x61, 0x72, 0x74, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x12, 0x1b, 0x2e, 0x6c,
0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x53, 0x74, 0x61, 0x72, 0x74, 0x45, 0x67, 0x72, 0x65,
0x73, 0x73, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x13, 0x2e, 0x6c, 0x69, 0x76, 0x65,
0x6b, 0x69, 0x74, 0x2e, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x6e, 0x66, 0x6f, 0x22, 0x06,
0xb2, 0x89, 0x01, 0x02, 0x20, 0x01, 0x12, 0x57, 0x0a, 0x10, 0x4c, 0x69, 0x73, 0x74, 0x41, 0x63,
0x74, 0x69, 0x76, 0x65, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x12, 0x1c, 0x2e, 0x72, 0x70, 0x63,
0x2e, 0x4c, 0x69, 0x73, 0x74, 0x41, 0x63, 0x74, 0x69, 0x76, 0x65, 0x45, 0x67, 0x72, 0x65, 0x73,
0x73, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x1d, 0x2e, 0x72, 0x70, 0x63, 0x2e, 0x4c,
0x69, 0x73, 0x74, 0x41, 0x63, 0x74, 0x69, 0x76, 0x65, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x52,
0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x22, 0x06, 0xb2, 0x89, 0x01, 0x02, 0x08, 0x01, 0x32,
0xa1, 0x01, 0x0a, 0x0d, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x48, 0x61, 0x6e, 0x64, 0x6c, 0x65,
0x72, 0x12, 0x49, 0x0a, 0x0c, 0x55, 0x70, 0x64, 0x61, 0x74, 0x65, 0x53, 0x74, 0x72, 0x65, 0x61,
0x6d, 0x12, 0x1c, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x55, 0x70, 0x64, 0x61,
0x74, 0x65, 0x53, 0x74, 0x72, 0x65, 0x61, 0x6d, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a,
0x13, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73,
0x49, 0x6e, 0x66, 0x6f, 0x22, 0x06, 0xb2, 0x89, 0x01, 0x02, 0x18, 0x01, 0x12, 0x35, 0x0a, 0x0a,
0x49, 0x6e, 0x66, 0x6f, 0x55, 0x70, 0x64, 0x61, 0x74, 0x65, 0x12, 0x0a, 0x2e, 0x72, 0x70, 0x63,
0x2e, 0x45, 0x6d, 0x70, 0x74, 0x79, 0x1a, 0x13, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74,
0x49, 0x6e, 0x66, 0x6f, 0x22, 0x06, 0xb2, 0x89, 0x01, 0x02, 0x18, 0x01, 0x12, 0x45, 0x0a, 0x0a,
0x53, 0x74, 0x6f, 0x70, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x12, 0x1a, 0x2e, 0x6c, 0x69, 0x76,
0x65, 0x6b, 0x69, 0x74, 0x2e, 0x53, 0x74, 0x6f, 0x70, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x52,
0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x13, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74,
0x2e, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x6e, 0x66, 0x6f, 0x22, 0x06, 0xb2, 0x89, 0x01,
0x02, 0x10, 0x01, 0x42, 0x2c, 0x5a, 0x2a, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f,
0x02, 0x18, 0x01, 0x42, 0x2c, 0x5a, 0x2a, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f,
0x6d, 0x2f, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2f, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69,
0x74, 0x2f, 0x70, 0x6b, 0x67, 0x2f, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x2f, 0x72, 0x70,
0x63, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33,
@@ -203,29 +161,26 @@ func file_pkg_service_rpc_egress_proto_rawDescGZIP() []byte {
return file_pkg_service_rpc_egress_proto_rawDescData
}
var file_pkg_service_rpc_egress_proto_msgTypes = make([]protoimpl.MessageInfo, 3)
var file_pkg_service_rpc_egress_proto_msgTypes = make([]protoimpl.MessageInfo, 2)
var file_pkg_service_rpc_egress_proto_goTypes = []interface{}{
(*Empty)(nil), // 0: rpc.Empty
(*ListActiveEgressRequest)(nil), // 1: rpc.ListActiveEgressRequest
(*ListActiveEgressResponse)(nil), // 2: rpc.ListActiveEgressResponse
(*livekit.StartEgressRequest)(nil), // 3: livekit.StartEgressRequest
(*livekit.UpdateStreamRequest)(nil), // 4: livekit.UpdateStreamRequest
(*livekit.StopEgressRequest)(nil), // 5: livekit.StopEgressRequest
(*livekit.EgressInfo)(nil), // 6: livekit.EgressInfo
(*ListActiveEgressRequest)(nil), // 0: rpc.ListActiveEgressRequest
(*ListActiveEgressResponse)(nil), // 1: rpc.ListActiveEgressResponse
(*livekit.StartEgressRequest)(nil), // 2: livekit.StartEgressRequest
(*livekit.UpdateStreamRequest)(nil), // 3: livekit.UpdateStreamRequest
(*livekit.StopEgressRequest)(nil), // 4: livekit.StopEgressRequest
(*livekit.EgressInfo)(nil), // 5: livekit.EgressInfo
}
var file_pkg_service_rpc_egress_proto_depIdxs = []int32{
3, // 0: rpc.EgressInternal.StartEgress:input_type -> livekit.StartEgressRequest
1, // 1: rpc.EgressInternal.ListActiveEgress:input_type -> rpc.ListActiveEgressRequest
4, // 2: rpc.EgressHandler.UpdateStream:input_type -> livekit.UpdateStreamRequest
5, // 3: rpc.EgressHandler.StopEgress:input_type -> livekit.StopEgressRequest
0, // 4: rpc.EgressHandler.InfoUpdate:input_type -> rpc.Empty
6, // 5: rpc.EgressInternal.StartEgress:output_type -> livekit.EgressInfo
2, // 6: rpc.EgressInternal.ListActiveEgress:output_type -> rpc.ListActiveEgressResponse
6, // 7: rpc.EgressHandler.UpdateStream:output_type -> livekit.EgressInfo
6, // 8: rpc.EgressHandler.StopEgress:output_type -> livekit.EgressInfo
6, // 9: rpc.EgressHandler.InfoUpdate:output_type -> livekit.EgressInfo
5, // [5:10] is the sub-list for method output_type
0, // [0:5] is the sub-list for method input_type
2, // 0: rpc.EgressInternal.StartEgress:input_type -> livekit.StartEgressRequest
0, // 1: rpc.EgressInternal.ListActiveEgress:input_type -> rpc.ListActiveEgressRequest
3, // 2: rpc.EgressHandler.UpdateStream:input_type -> livekit.UpdateStreamRequest
4, // 3: rpc.EgressHandler.StopEgress:input_type -> livekit.StopEgressRequest
5, // 4: rpc.EgressInternal.StartEgress:output_type -> livekit.EgressInfo
1, // 5: rpc.EgressInternal.ListActiveEgress:output_type -> rpc.ListActiveEgressResponse
5, // 6: rpc.EgressHandler.UpdateStream:output_type -> livekit.EgressInfo
5, // 7: rpc.EgressHandler.StopEgress:output_type -> livekit.EgressInfo
4, // [4:8] is the sub-list for method output_type
0, // [0:4] is the sub-list for method input_type
0, // [0:0] is the sub-list for extension type_name
0, // [0:0] is the sub-list for extension extendee
0, // [0:0] is the sub-list for field type_name
@@ -238,18 +193,6 @@ func file_pkg_service_rpc_egress_proto_init() {
}
if !protoimpl.UnsafeEnabled {
file_pkg_service_rpc_egress_proto_msgTypes[0].Exporter = func(v interface{}, i int) interface{} {
switch v := v.(*Empty); i {
case 0:
return &v.state
case 1:
return &v.sizeCache
case 2:
return &v.unknownFields
default:
return nil
}
}
file_pkg_service_rpc_egress_proto_msgTypes[1].Exporter = func(v interface{}, i int) interface{} {
switch v := v.(*ListActiveEgressRequest); i {
case 0:
return &v.state
@@ -261,7 +204,7 @@ func file_pkg_service_rpc_egress_proto_init() {
return nil
}
}
file_pkg_service_rpc_egress_proto_msgTypes[2].Exporter = func(v interface{}, i int) interface{} {
file_pkg_service_rpc_egress_proto_msgTypes[1].Exporter = func(v interface{}, i int) interface{} {
switch v := v.(*ListActiveEgressResponse); i {
case 0:
return &v.state
@@ -280,7 +223,7 @@ func file_pkg_service_rpc_egress_proto_init() {
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
RawDescriptor: file_pkg_service_rpc_egress_proto_rawDesc,
NumEnums: 0,
NumMessages: 3,
NumMessages: 2,
NumExtensions: 0,
NumServices: 2,
},
-5
View File
@@ -24,13 +24,8 @@ service EgressHandler {
rpc StopEgress(livekit.StopEgressRequest) returns (livekit.EgressInfo) {
option (psrpc.options).topics = true;
}
rpc InfoUpdate(Empty) returns (livekit.EgressInfo) {
option (psrpc.options).subscription = true;
}
}
message Empty {}
message ListActiveEgressRequest {}
message ListActiveEgressResponse {
+22 -35
View File
@@ -1,4 +1,4 @@
// Code generated by protoc-gen-psrpc v0.2.2, DO NOT EDIT.
// Code generated by protoc-gen-psrpc v0.2.3, DO NOT EDIT.
// source: pkg/service/rpc/egress.proto
package rpc
@@ -119,8 +119,6 @@ type EgressHandlerClient interface {
UpdateStream(context.Context, string, *livekit.UpdateStreamRequest, ...psrpc1.RequestOption) (*livekit.EgressInfo, error)
StopEgress(context.Context, string, *livekit.StopEgressRequest, ...psrpc1.RequestOption) (*livekit.EgressInfo, error)
SubscribeInfoUpdate(context.Context) (psrpc1.Subscription[*livekit.EgressInfo], error)
}
// ==================================
@@ -144,8 +142,6 @@ type EgressHandlerServer interface {
RegisterStopEgressTopic(string) error
DeregisterStopEgressTopic(string)
PublishInfoUpdate(context.Context, *livekit.EgressInfo) error
// Close and wait for pending RPCs to complete
Shutdown()
@@ -181,10 +177,6 @@ func (c *egressHandlerClient) StopEgress(ctx context.Context, topic string, req
return psrpc1.RequestSingle[*livekit.EgressInfo](ctx, c.client, "StopEgress", topic, req, opts...)
}
func (c *egressHandlerClient) SubscribeInfoUpdate(ctx context.Context) (psrpc1.Subscription[*livekit.EgressInfo], error) {
return psrpc1.JoinQueue[*livekit.EgressInfo](ctx, c.client, "InfoUpdate", "")
}
// ====================
// EgressHandler Server
// ====================
@@ -221,10 +213,6 @@ func (s *egressHandlerServer) DeregisterStopEgressTopic(topic string) {
s.rpc.DeregisterHandler("StopEgress", topic)
}
func (s *egressHandlerServer) PublishInfoUpdate(ctx context.Context, msg *livekit.EgressInfo) error {
return s.rpc.Publish(ctx, "InfoUpdate", "", msg)
}
func (s *egressHandlerServer) Shutdown() {
s.rpc.Close(false)
}
@@ -234,26 +222,25 @@ func (s *egressHandlerServer) Kill() {
}
var psrpcFileDescriptor0 = []byte{
// 329 bytes of a gzipped FileDescriptorProto
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x8c, 0x52, 0xd1, 0x4a, 0x02, 0x41,
0x14, 0x65, 0x92, 0x2c, 0x6f, 0x19, 0x32, 0x05, 0xd9, 0xa6, 0x20, 0xfb, 0x14, 0x11, 0xbb, 0x60,
0xf4, 0xd0, 0x63, 0x81, 0x94, 0xd0, 0x93, 0x12, 0x41, 0x2f, 0xb2, 0xce, 0xde, 0x6c, 0x50, 0x77,
0xa6, 0x99, 0xab, 0xd0, 0x27, 0xf4, 0x3b, 0x7e, 0x4d, 0x9f, 0x13, 0xed, 0x8c, 0xb2, 0x18, 0x56,
0x4f, 0xc3, 0x9c, 0x73, 0xef, 0xb9, 0xe7, 0x70, 0x2f, 0x34, 0xf4, 0x78, 0x14, 0x5b, 0x34, 0x73,
0x29, 0x30, 0x36, 0x5a, 0xc4, 0x38, 0x32, 0x68, 0x6d, 0xa4, 0x8d, 0x22, 0xc5, 0x4b, 0x46, 0x8b,
0xa0, 0xaa, 0x34, 0x49, 0x95, 0x79, 0x2c, 0x08, 0x26, 0x72, 0x8e, 0x63, 0x49, 0x03, 0xa3, 0xc5,
0x40, 0x66, 0x84, 0x26, 0x4b, 0x26, 0x9e, 0x3b, 0x5a, 0x72, 0x45, 0x95, 0x70, 0x07, 0xb6, 0x3b,
0x53, 0x4d, 0xef, 0xe1, 0x09, 0x1c, 0x3f, 0x48, 0x4b, 0x37, 0x82, 0xe4, 0x1c, 0x3b, 0x79, 0x49,
0x0f, 0xdf, 0x66, 0x68, 0x29, 0xbc, 0x86, 0xfa, 0x4f, 0xca, 0x6a, 0x95, 0x59, 0xe4, 0x4d, 0x00,
0xa7, 0x37, 0x90, 0xa9, 0xad, 0xb3, 0x56, 0xe9, 0xac, 0xd2, 0xab, 0x38, 0xa4, 0x9b, 0xda, 0xf6,
0x82, 0xc1, 0x81, 0xeb, 0xe8, 0x7a, 0x37, 0xfc, 0x0e, 0xf6, 0xfa, 0x94, 0x18, 0x72, 0x30, 0x3f,
0x8d, 0xbc, 0xaf, 0xa8, 0x80, 0xfa, 0xc9, 0xc1, 0xe1, 0x8a, 0x5c, 0x8a, 0xbc, 0xa8, 0xb0, 0xbc,
0xf8, 0x60, 0x5b, 0x2d, 0xc6, 0x9f, 0xa0, 0xb6, 0x6e, 0x8b, 0x37, 0x22, 0xa3, 0x45, 0xb4, 0x21,
0x48, 0xd0, 0xdc, 0xc0, 0xba, 0x2c, 0x4e, 0x78, 0x97, 0xb5, 0x3f, 0x19, 0x54, 0x1d, 0x75, 0x9f,
0x64, 0xe9, 0x04, 0x0d, 0xef, 0xc2, 0xfe, 0xa3, 0x4e, 0x13, 0xc2, 0x3e, 0x19, 0x4c, 0xa6, 0xbc,
0xb1, 0xf2, 0x55, 0x84, 0xff, 0x76, 0x5d, 0x67, 0xbc, 0x03, 0xd0, 0x27, 0xa5, 0xbd, 0xdf, 0xa0,
0x90, 0x7e, 0x09, 0xfe, 0x4b, 0xe6, 0x0a, 0xe0, 0xfb, 0xef, 0xc6, 0x73, 0xc8, 0x83, 0xe5, 0x8b,
0xfc, 0xa5, 0xad, 0xc6, 0x6e, 0x2f, 0x9e, 0xcf, 0x47, 0x92, 0x5e, 0x67, 0xc3, 0x48, 0xa8, 0x69,
0xec, 0x0b, 0x57, 0xef, 0xda, 0xbd, 0x0d, 0xcb, 0xf9, 0x8d, 0x5c, 0x7e, 0x05, 0x00, 0x00, 0xff,
0xff, 0x26, 0x11, 0xba, 0xff, 0x89, 0x02, 0x00, 0x00,
// 305 bytes of a gzipped FileDescriptorProto
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x8c, 0x92, 0xc1, 0x4a, 0x2b, 0x31,
0x14, 0x86, 0xc9, 0x2d, 0x94, 0xdb, 0xa3, 0x15, 0x89, 0x82, 0x35, 0xb6, 0x50, 0xba, 0x12, 0x91,
0x0c, 0xd4, 0x95, 0x4b, 0x85, 0xa2, 0x05, 0x57, 0x2d, 0x22, 0xb8, 0x29, 0xd3, 0xcc, 0xb1, 0x86,
0xb6, 0x93, 0x98, 0x9c, 0xf6, 0x1d, 0x7c, 0x0c, 0x5f, 0xa1, 0x4f, 0x28, 0x34, 0x99, 0x61, 0xa8,
0x14, 0x5d, 0x05, 0xfe, 0x2f, 0xf9, 0xf3, 0x1d, 0x12, 0x68, 0xdb, 0xf9, 0x2c, 0xf1, 0xe8, 0xd6,
0x5a, 0x61, 0xe2, 0xac, 0x4a, 0x70, 0xe6, 0xd0, 0x7b, 0x69, 0x9d, 0x21, 0xc3, 0x6b, 0xce, 0x2a,
0xd1, 0x34, 0x96, 0xb4, 0xc9, 0x63, 0x26, 0xc4, 0x42, 0xaf, 0x71, 0xae, 0x69, 0xe2, 0xac, 0x9a,
0xe8, 0x9c, 0xd0, 0xe5, 0xe9, 0x22, 0xb2, 0xd3, 0x82, 0x55, 0x5b, 0x7a, 0xe7, 0x70, 0xf6, 0xa4,
0x3d, 0xdd, 0x29, 0xd2, 0x6b, 0x1c, 0x6c, 0xc9, 0x08, 0x3f, 0x56, 0xe8, 0xa9, 0x77, 0x0b, 0xad,
0x9f, 0xc8, 0x5b, 0x93, 0x7b, 0xe4, 0x1d, 0x80, 0x50, 0x33, 0xd1, 0x99, 0x6f, 0xb1, 0x6e, 0xed,
0xb2, 0x31, 0x6a, 0x84, 0x64, 0x98, 0xf9, 0xfe, 0x86, 0xc1, 0x51, 0x38, 0x31, 0x8c, 0x12, 0xfc,
0x01, 0x0e, 0xc6, 0x94, 0x3a, 0x0a, 0x31, 0xbf, 0x90, 0x51, 0x47, 0x56, 0xd2, 0x78, 0xb3, 0x38,
0x29, 0x61, 0x51, 0xf2, 0x66, 0x7a, 0xf5, 0xcd, 0x27, 0xfb, 0xd7, 0x65, 0xfc, 0x05, 0x8e, 0x77,
0xb5, 0x78, 0x5b, 0x3a, 0xab, 0xe4, 0x9e, 0x41, 0x44, 0x67, 0x0f, 0x0d, 0xb3, 0x84, 0xe2, 0xff,
0xac, 0xff, 0xc5, 0xa0, 0x19, 0xd0, 0x63, 0x9a, 0x67, 0x0b, 0x74, 0x7c, 0x08, 0x87, 0xcf, 0x36,
0x4b, 0x09, 0xc7, 0xe4, 0x30, 0x5d, 0xf2, 0x76, 0xe9, 0x55, 0x8d, 0x7f, 0xb7, 0x6e, 0x31, 0x3e,
0x00, 0x18, 0x93, 0xb1, 0xd1, 0x57, 0x54, 0xa6, 0x2f, 0xc2, 0xbf, 0xd4, 0xdc, 0x5f, 0xbf, 0x5e,
0xcd, 0x34, 0xbd, 0xaf, 0xa6, 0x52, 0x99, 0x65, 0x12, 0x37, 0x96, 0xeb, 0xce, 0x7f, 0x99, 0xd6,
0xb7, 0x6f, 0x7c, 0xf3, 0x1d, 0x00, 0x00, 0xff, 0xff, 0x58, 0x8f, 0x8d, 0xdd, 0x49, 0x02, 0x00,
0x00,
}
+2 -2
View File
@@ -1,7 +1,7 @@
// Code generated by protoc-gen-go. DO NOT EDIT.
// versions:
// protoc-gen-go v1.28.1
// protoc v3.21.12
// protoc-gen-go v1.27.1
// protoc v3.21.5
// source: pkg/service/rpc/ingress.proto
package rpc
+1 -1
View File
@@ -1,4 +1,4 @@
// Code generated by protoc-gen-psrpc v0.2.2, DO NOT EDIT.
// Code generated by protoc-gen-psrpc v0.2.3, DO NOT EDIT.
// source: pkg/service/rpc/ingress.proto
package rpc
+39 -95
View File
@@ -1,7 +1,7 @@
// Code generated by protoc-gen-go. DO NOT EDIT.
// versions:
// protoc-gen-go v1.28.1
// protoc v3.21.12
// protoc-gen-go v1.27.1
// protoc v3.21.5
// source: pkg/service/rpc/io.proto
package rpc
@@ -10,8 +10,8 @@ import (
livekit "github.com/livekit/protocol/livekit"
protoreflect "google.golang.org/protobuf/reflect/protoreflect"
protoimpl "google.golang.org/protobuf/runtime/protoimpl"
emptypb "google.golang.org/protobuf/types/known/emptypb"
reflect "reflect"
sync "sync"
)
const (
@@ -21,94 +21,53 @@ const (
_ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20)
)
type Ignored struct {
state protoimpl.MessageState
sizeCache protoimpl.SizeCache
unknownFields protoimpl.UnknownFields
}
func (x *Ignored) Reset() {
*x = Ignored{}
if protoimpl.UnsafeEnabled {
mi := &file_pkg_service_rpc_io_proto_msgTypes[0]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
}
func (x *Ignored) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*Ignored) ProtoMessage() {}
func (x *Ignored) ProtoReflect() protoreflect.Message {
mi := &file_pkg_service_rpc_io_proto_msgTypes[0]
if protoimpl.UnsafeEnabled && x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use Ignored.ProtoReflect.Descriptor instead.
func (*Ignored) Descriptor() ([]byte, []int) {
return file_pkg_service_rpc_io_proto_rawDescGZIP(), []int{0}
}
var File_pkg_service_rpc_io_proto protoreflect.FileDescriptor
var file_pkg_service_rpc_io_proto_rawDesc = []byte{
0x0a, 0x18, 0x70, 0x6b, 0x67, 0x2f, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x2f, 0x72, 0x70,
0x63, 0x2f, 0x69, 0x6f, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x12, 0x03, 0x72, 0x70, 0x63, 0x1a,
0x1a, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x5f, 0x72, 0x70, 0x63, 0x5f, 0x69, 0x6e, 0x74,
0x65, 0x72, 0x6e, 0x61, 0x6c, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x22, 0x09, 0x0a, 0x07, 0x49,
0x67, 0x6e, 0x6f, 0x72, 0x65, 0x64, 0x32, 0xa3, 0x01, 0x0a, 0x06, 0x49, 0x4f, 0x49, 0x6e, 0x66,
0x6f, 0x12, 0x51, 0x0a, 0x0e, 0x47, 0x65, 0x74, 0x49, 0x6e, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49,
0x6e, 0x66, 0x6f, 0x12, 0x1e, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x47, 0x65,
0x74, 0x49, 0x6e, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x6e, 0x66, 0x6f, 0x52, 0x65, 0x71, 0x75,
0x65, 0x73, 0x74, 0x1a, 0x1f, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x47, 0x65,
0x74, 0x49, 0x6e, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x6e, 0x66, 0x6f, 0x52, 0x65, 0x73, 0x70,
0x6f, 0x6e, 0x73, 0x65, 0x12, 0x46, 0x0a, 0x12, 0x55, 0x70, 0x64, 0x61, 0x74, 0x65, 0x49, 0x6e,
0x67, 0x72, 0x65, 0x73, 0x73, 0x53, 0x74, 0x61, 0x74, 0x65, 0x12, 0x22, 0x2e, 0x6c, 0x69, 0x76,
0x65, 0x6b, 0x69, 0x74, 0x2e, 0x55, 0x70, 0x64, 0x61, 0x74, 0x65, 0x49, 0x6e, 0x67, 0x72, 0x65,
0x73, 0x73, 0x53, 0x74, 0x61, 0x74, 0x65, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x0c,
0x2e, 0x72, 0x70, 0x63, 0x2e, 0x49, 0x67, 0x6e, 0x6f, 0x72, 0x65, 0x64, 0x42, 0x2c, 0x5a, 0x2a,
0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x6c, 0x69, 0x76, 0x65, 0x6b,
0x69, 0x74, 0x2f, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2f, 0x70, 0x6b, 0x67, 0x2f, 0x73,
0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x2f, 0x72, 0x70, 0x63, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74,
0x6f, 0x33,
0x14, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x5f, 0x65, 0x67, 0x72, 0x65, 0x73, 0x73, 0x2e,
0x70, 0x72, 0x6f, 0x74, 0x6f, 0x1a, 0x1a, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x5f, 0x72,
0x70, 0x63, 0x5f, 0x69, 0x6e, 0x74, 0x65, 0x72, 0x6e, 0x61, 0x6c, 0x2e, 0x70, 0x72, 0x6f, 0x74,
0x6f, 0x1a, 0x1b, 0x67, 0x6f, 0x6f, 0x67, 0x6c, 0x65, 0x2f, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x62,
0x75, 0x66, 0x2f, 0x65, 0x6d, 0x70, 0x74, 0x79, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x32, 0xee,
0x01, 0x0a, 0x06, 0x49, 0x4f, 0x49, 0x6e, 0x66, 0x6f, 0x12, 0x3f, 0x0a, 0x10, 0x55, 0x70, 0x64,
0x61, 0x74, 0x65, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x6e, 0x66, 0x6f, 0x12, 0x13, 0x2e,
0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x45, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x6e,
0x66, 0x6f, 0x1a, 0x16, 0x2e, 0x67, 0x6f, 0x6f, 0x67, 0x6c, 0x65, 0x2e, 0x70, 0x72, 0x6f, 0x74,
0x6f, 0x62, 0x75, 0x66, 0x2e, 0x45, 0x6d, 0x70, 0x74, 0x79, 0x12, 0x51, 0x0a, 0x0e, 0x47, 0x65,
0x74, 0x49, 0x6e, 0x67, 0x72, 0x65, 0x73, 0x73, 0x49, 0x6e, 0x66, 0x6f, 0x12, 0x1e, 0x2e, 0x6c,
0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x47, 0x65, 0x74, 0x49, 0x6e, 0x67, 0x72, 0x65, 0x73,
0x73, 0x49, 0x6e, 0x66, 0x6f, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x1f, 0x2e, 0x6c,
0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x47, 0x65, 0x74, 0x49, 0x6e, 0x67, 0x72, 0x65, 0x73,
0x73, 0x49, 0x6e, 0x66, 0x6f, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x50, 0x0a,
0x12, 0x55, 0x70, 0x64, 0x61, 0x74, 0x65, 0x49, 0x6e, 0x67, 0x72, 0x65, 0x73, 0x73, 0x53, 0x74,
0x61, 0x74, 0x65, 0x12, 0x22, 0x2e, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2e, 0x55, 0x70,
0x64, 0x61, 0x74, 0x65, 0x49, 0x6e, 0x67, 0x72, 0x65, 0x73, 0x73, 0x53, 0x74, 0x61, 0x74, 0x65,
0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x16, 0x2e, 0x67, 0x6f, 0x6f, 0x67, 0x6c, 0x65,
0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x62, 0x75, 0x66, 0x2e, 0x45, 0x6d, 0x70, 0x74, 0x79, 0x42,
0x2c, 0x5a, 0x2a, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x6c, 0x69,
0x76, 0x65, 0x6b, 0x69, 0x74, 0x2f, 0x6c, 0x69, 0x76, 0x65, 0x6b, 0x69, 0x74, 0x2f, 0x70, 0x6b,
0x67, 0x2f, 0x73, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x2f, 0x72, 0x70, 0x63, 0x62, 0x06, 0x70,
0x72, 0x6f, 0x74, 0x6f, 0x33,
}
var (
file_pkg_service_rpc_io_proto_rawDescOnce sync.Once
file_pkg_service_rpc_io_proto_rawDescData = file_pkg_service_rpc_io_proto_rawDesc
)
func file_pkg_service_rpc_io_proto_rawDescGZIP() []byte {
file_pkg_service_rpc_io_proto_rawDescOnce.Do(func() {
file_pkg_service_rpc_io_proto_rawDescData = protoimpl.X.CompressGZIP(file_pkg_service_rpc_io_proto_rawDescData)
})
return file_pkg_service_rpc_io_proto_rawDescData
}
var file_pkg_service_rpc_io_proto_msgTypes = make([]protoimpl.MessageInfo, 1)
var file_pkg_service_rpc_io_proto_goTypes = []interface{}{
(*Ignored)(nil), // 0: rpc.Ignored
(*livekit.EgressInfo)(nil), // 0: livekit.EgressInfo
(*livekit.GetIngressInfoRequest)(nil), // 1: livekit.GetIngressInfoRequest
(*livekit.UpdateIngressStateRequest)(nil), // 2: livekit.UpdateIngressStateRequest
(*livekit.GetIngressInfoResponse)(nil), // 3: livekit.GetIngressInfoResponse
(*emptypb.Empty)(nil), // 3: google.protobuf.Empty
(*livekit.GetIngressInfoResponse)(nil), // 4: livekit.GetIngressInfoResponse
}
var file_pkg_service_rpc_io_proto_depIdxs = []int32{
1, // 0: rpc.IOInfo.GetIngressInfo:input_type -> livekit.GetIngressInfoRequest
2, // 1: rpc.IOInfo.UpdateIngressState:input_type -> livekit.UpdateIngressStateRequest
3, // 2: rpc.IOInfo.GetIngressInfo:output_type -> livekit.GetIngressInfoResponse
0, // 3: rpc.IOInfo.UpdateIngressState:output_type -> rpc.Ignored
2, // [2:4] is the sub-list for method output_type
0, // [0:2] is the sub-list for method input_type
0, // 0: rpc.IOInfo.UpdateEgressInfo:input_type -> livekit.EgressInfo
1, // 1: rpc.IOInfo.GetIngressInfo:input_type -> livekit.GetIngressInfoRequest
2, // 2: rpc.IOInfo.UpdateIngressState:input_type -> livekit.UpdateIngressStateRequest
3, // 3: rpc.IOInfo.UpdateEgressInfo:output_type -> google.protobuf.Empty
4, // 4: rpc.IOInfo.GetIngressInfo:output_type -> livekit.GetIngressInfoResponse
3, // 5: rpc.IOInfo.UpdateIngressState:output_type -> google.protobuf.Empty
3, // [3:6] is the sub-list for method output_type
0, // [0:3] is the sub-list for method input_type
0, // [0:0] is the sub-list for extension type_name
0, // [0:0] is the sub-list for extension extendee
0, // [0:0] is the sub-list for field type_name
@@ -119,33 +78,18 @@ func file_pkg_service_rpc_io_proto_init() {
if File_pkg_service_rpc_io_proto != nil {
return
}
if !protoimpl.UnsafeEnabled {
file_pkg_service_rpc_io_proto_msgTypes[0].Exporter = func(v interface{}, i int) interface{} {
switch v := v.(*Ignored); i {
case 0:
return &v.state
case 1:
return &v.sizeCache
case 2:
return &v.unknownFields
default:
return nil
}
}
}
type x struct{}
out := protoimpl.TypeBuilder{
File: protoimpl.DescBuilder{
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
RawDescriptor: file_pkg_service_rpc_io_proto_rawDesc,
NumEnums: 0,
NumMessages: 1,
NumMessages: 0,
NumExtensions: 0,
NumServices: 1,
},
GoTypes: file_pkg_service_rpc_io_proto_goTypes,
DependencyIndexes: file_pkg_service_rpc_io_proto_depIdxs,
MessageInfos: file_pkg_service_rpc_io_proto_msgTypes,
}.Build()
File_pkg_service_rpc_io_proto = out.File
file_pkg_service_rpc_io_proto_rawDesc = nil
+4 -3
View File
@@ -4,11 +4,12 @@ package rpc;
option go_package = "github.com/livekit/livekit/pkg/service/rpc";
import "livekit_egress.proto";
import "livekit_rpc_internal.proto";
import "google/protobuf/empty.proto";
service IOInfo {
rpc UpdateEgressInfo(livekit.EgressInfo) returns (google.protobuf.Empty);
rpc GetIngressInfo(livekit.GetIngressInfoRequest) returns (livekit.GetIngressInfoResponse);
rpc UpdateIngressState(livekit.UpdateIngressStateRequest) returns (Ignored);
rpc UpdateIngressState(livekit.UpdateIngressStateRequest) returns (google.protobuf.Empty);
}
message Ignored {}
+37 -19
View File
@@ -1,4 +1,4 @@
// Code generated by protoc-gen-psrpc v0.2.2, DO NOT EDIT.
// Code generated by protoc-gen-psrpc v0.2.3, DO NOT EDIT.
// source: pkg/service/rpc/io.proto
package rpc
@@ -6,6 +6,8 @@ package rpc
import context "context"
import psrpc1 "github.com/livekit/psrpc"
import google_protobuf2 "google.golang.org/protobuf/types/known/emptypb"
import livekit "github.com/livekit/protocol/livekit"
import livekit3 "github.com/livekit/protocol/livekit"
// =======================
@@ -13,9 +15,11 @@ import livekit3 "github.com/livekit/protocol/livekit"
// =======================
type IOInfoClient interface {
UpdateEgressInfo(context.Context, *livekit.EgressInfo, ...psrpc1.RequestOption) (*google_protobuf2.Empty, error)
GetIngressInfo(context.Context, *livekit3.GetIngressInfoRequest, ...psrpc1.RequestOption) (*livekit3.GetIngressInfoResponse, error)
UpdateIngressState(context.Context, *livekit3.UpdateIngressStateRequest, ...psrpc1.RequestOption) (*Ignored, error)
UpdateIngressState(context.Context, *livekit3.UpdateIngressStateRequest, ...psrpc1.RequestOption) (*google_protobuf2.Empty, error)
}
// ===========================
@@ -23,9 +27,11 @@ type IOInfoClient interface {
// ===========================
type IOInfoServerImpl interface {
UpdateEgressInfo(context.Context, *livekit.EgressInfo) (*google_protobuf2.Empty, error)
GetIngressInfo(context.Context, *livekit3.GetIngressInfoRequest) (*livekit3.GetIngressInfoResponse, error)
UpdateIngressState(context.Context, *livekit3.UpdateIngressStateRequest) (*Ignored, error)
UpdateIngressState(context.Context, *livekit3.UpdateIngressStateRequest) (*google_protobuf2.Empty, error)
}
// =======================
@@ -60,12 +66,16 @@ func NewIOInfoClient(clientID string, bus psrpc1.MessageBus, opts ...psrpc1.Clie
}, nil
}
func (c *iOInfoClient) UpdateEgressInfo(ctx context.Context, req *livekit.EgressInfo, opts ...psrpc1.RequestOption) (*google_protobuf2.Empty, error) {
return psrpc1.RequestSingle[*google_protobuf2.Empty](ctx, c.client, "UpdateEgressInfo", "", req, opts...)
}
func (c *iOInfoClient) GetIngressInfo(ctx context.Context, req *livekit3.GetIngressInfoRequest, opts ...psrpc1.RequestOption) (*livekit3.GetIngressInfoResponse, error) {
return psrpc1.RequestSingle[*livekit3.GetIngressInfoResponse](ctx, c.client, "GetIngressInfo", "", req, opts...)
}
func (c *iOInfoClient) UpdateIngressState(ctx context.Context, req *livekit3.UpdateIngressStateRequest, opts ...psrpc1.RequestOption) (*Ignored, error) {
return psrpc1.RequestSingle[*Ignored](ctx, c.client, "UpdateIngressState", "", req, opts...)
func (c *iOInfoClient) UpdateIngressState(ctx context.Context, req *livekit3.UpdateIngressStateRequest, opts ...psrpc1.RequestOption) (*google_protobuf2.Empty, error) {
return psrpc1.RequestSingle[*google_protobuf2.Empty](ctx, c.client, "UpdateIngressState", "", req, opts...)
}
// =============
@@ -83,6 +93,12 @@ func NewIOInfoServer(serverID string, svc IOInfoServerImpl, bus psrpc1.MessageBu
s := psrpc1.NewRPCServer("IOInfo", serverID, bus, opts...)
var err error
err = psrpc1.RegisterHandler(s, "UpdateEgressInfo", "", svc.UpdateEgressInfo, nil)
if err != nil {
s.Close(false)
return nil, err
}
err = psrpc1.RegisterHandler(s, "GetIngressInfo", "", svc.GetIngressInfo, nil)
if err != nil {
s.Close(false)
@@ -110,18 +126,20 @@ func (s *iOInfoServer) Kill() {
}
var psrpcFileDescriptor2 = []byte{
// 201 bytes of a gzipped FileDescriptorProto
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0xe2, 0x92, 0x28, 0xc8, 0x4e, 0xd7,
0x2f, 0x4e, 0x2d, 0x2a, 0xcb, 0x4c, 0x4e, 0xd5, 0x2f, 0x2a, 0x48, 0xd6, 0xcf, 0xcc, 0xd7, 0x2b,
0x28, 0xca, 0x2f, 0xc9, 0x17, 0x62, 0x2e, 0x2a, 0x48, 0x96, 0x92, 0xca, 0xc9, 0x2c, 0x4b, 0xcd,
0xce, 0x2c, 0x89, 0x2f, 0x2a, 0x48, 0x8e, 0xcf, 0xcc, 0x2b, 0x49, 0x2d, 0xca, 0x4b, 0xcc, 0x81,
0x28, 0x50, 0xe2, 0xe4, 0x62, 0xf7, 0x4c, 0xcf, 0xcb, 0x2f, 0x4a, 0x4d, 0x31, 0x5a, 0xcc, 0xc8,
0xc5, 0xe6, 0xe9, 0xef, 0x99, 0x97, 0x96, 0x2f, 0x14, 0xc8, 0xc5, 0xe7, 0x9e, 0x5a, 0xe2, 0x99,
0x97, 0x5e, 0x94, 0x5a, 0x5c, 0x0c, 0x16, 0x91, 0xd3, 0x83, 0x1a, 0xa2, 0x87, 0x2a, 0x11, 0x94,
0x5a, 0x58, 0x9a, 0x5a, 0x5c, 0x22, 0x25, 0x8f, 0x53, 0xbe, 0xb8, 0x20, 0x3f, 0xaf, 0x38, 0x55,
0xc8, 0x8d, 0x4b, 0x28, 0xb4, 0x20, 0x25, 0xb1, 0x24, 0x15, 0x2a, 0x19, 0x5c, 0x92, 0x58, 0x92,
0x2a, 0xa4, 0x04, 0xd7, 0x86, 0x29, 0x09, 0x33, 0x9a, 0x47, 0xaf, 0xa8, 0x20, 0x59, 0x0f, 0xea,
0x4a, 0x27, 0x9d, 0x28, 0xad, 0xf4, 0xcc, 0x92, 0x8c, 0xd2, 0x24, 0xbd, 0xe4, 0xfc, 0x5c, 0x7d,
0xa8, 0x6e, 0x38, 0x8d, 0x16, 0x10, 0x49, 0x6c, 0x60, 0x5f, 0x1a, 0x03, 0x02, 0x00, 0x00, 0xff,
0xff, 0xba, 0x2f, 0x5a, 0x5b, 0x22, 0x01, 0x00, 0x00,
// 236 bytes of a gzipped FileDescriptorProto
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x74, 0x90, 0xcd, 0x4a, 0xc4, 0x30,
0x14, 0x85, 0x11, 0x61, 0x16, 0x59, 0x88, 0x44, 0x11, 0x89, 0xa0, 0xe0, 0x52, 0x24, 0x01, 0x7d,
0x00, 0x41, 0x18, 0xa4, 0x2b, 0xff, 0x70, 0xe3, 0x66, 0x68, 0xe3, 0x9d, 0x18, 0xa6, 0x93, 0x7b,
0x4d, 0x6e, 0x07, 0x7c, 0x69, 0x9f, 0x41, 0x6c, 0xd2, 0x0e, 0x2a, 0x5d, 0x85, 0x9c, 0x8f, 0xf3,
0x71, 0xb8, 0xe2, 0x98, 0x56, 0xce, 0x24, 0x88, 0x1b, 0x6f, 0xc1, 0x44, 0xb2, 0xc6, 0xa3, 0xa6,
0x88, 0x8c, 0x72, 0x37, 0x92, 0x55, 0x87, 0xad, 0xdf, 0xc0, 0xca, 0xf3, 0x02, 0x5c, 0x84, 0x94,
0x32, 0x52, 0x6a, 0x48, 0x23, 0xd9, 0x85, 0x0f, 0x0c, 0x31, 0xd4, 0x6d, 0x61, 0x27, 0x0e, 0xd1,
0xb5, 0x60, 0xfa, 0x5f, 0xd3, 0x2d, 0x0d, 0xac, 0x89, 0x3f, 0x33, 0xbc, 0xfa, 0xda, 0x11, 0xb3,
0xea, 0xbe, 0x0a, 0x4b, 0x94, 0x37, 0x62, 0xff, 0x85, 0xde, 0x6a, 0x86, 0x79, 0x6f, 0xee, 0xb3,
0x03, 0x5d, 0xc4, 0x7a, 0x1b, 0xaa, 0x23, 0x9d, 0x8d, 0x7a, 0x30, 0xea, 0xf9, 0x8f, 0x51, 0x3e,
0x8a, 0xbd, 0x3b, 0xe0, 0x2a, 0x6c, 0xeb, 0xa7, 0x63, 0xfd, 0x37, 0x78, 0x82, 0x8f, 0x0e, 0x12,
0xab, 0xb3, 0x49, 0x9e, 0x08, 0x43, 0x02, 0xf9, 0x20, 0x64, 0xde, 0x54, 0xe0, 0x33, 0xd7, 0x0c,
0xf2, 0x7c, 0xac, 0xfd, 0x87, 0x83, 0x7a, 0x62, 0xe4, 0xed, 0xe5, 0xeb, 0x85, 0xf3, 0xfc, 0xde,
0x35, 0xda, 0xe2, 0xda, 0x14, 0xcf, 0xf8, 0xfe, 0xb9, 0x7d, 0x33, 0xeb, 0xdb, 0xd7, 0xdf, 0x01,
0x00, 0x00, 0xff, 0xff, 0x87, 0x58, 0x4d, 0x3b, 0x95, 0x01, 0x00, 0x00,
}
+20 -24
View File
@@ -27,25 +27,25 @@ import (
)
type LivekitServer struct {
config *config.Config
egressService *EgressService
ingressService *IngressService
rtcService *RTCService
httpServer *http.Server
promServer *http.Server
router routing.Router
roomManager *RoomManager
turnServer *turn.Server
currentNode routing.LocalNode
running atomic.Bool
doneChan chan struct{}
closedChan chan struct{}
config *config.Config
ioService *IOInfoService
rtcService *RTCService
httpServer *http.Server
promServer *http.Server
router routing.Router
roomManager *RoomManager
turnServer *turn.Server
currentNode routing.LocalNode
running atomic.Bool
doneChan chan struct{}
closedChan chan struct{}
}
func NewLivekitServer(conf *config.Config,
roomService livekit.RoomService,
egressService *EgressService,
ingressService *IngressService,
ioService *IOInfoService,
rtcService *RTCService,
keyProvider auth.KeyProvider,
router routing.Router,
@@ -54,12 +54,11 @@ func NewLivekitServer(conf *config.Config,
currentNode routing.LocalNode,
) (s *LivekitServer, err error) {
s = &LivekitServer{
config: conf,
egressService: egressService,
ingressService: ingressService,
rtcService: rtcService,
router: router,
roomManager: roomManager,
config: conf,
ioService: ioService,
rtcService: rtcService,
router: router,
roomManager: roomManager,
// turn server starts automatically
turnServer: turnServer,
currentNode: currentNode,
@@ -154,12 +153,10 @@ func (s *LivekitServer) Start() error {
return err
}
if err := s.egressService.Start(); err != nil {
if err := s.ioService.Start(); err != nil {
return err
}
s.ingressService.Start()
addresses := s.config.BindAddresses
if addresses == nil {
addresses = []string{""}
@@ -248,8 +245,7 @@ func (s *LivekitServer) Start() error {
}
s.roomManager.Stop()
s.egressService.Stop()
s.ingressService.Stop()
s.ioService.Stop()
close(s.closedChan)
return nil
+1
View File
@@ -46,6 +46,7 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live
telemetry.NewAnalyticsService,
telemetry.NewTelemetryService,
getMessageBus,
NewIOInfoService,
getEgressClient,
egress.NewRedisRPCClient,
getEgressStore,
+5 -1
View File
@@ -80,6 +80,10 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live
ingressRPCClient := getIngressRPCClient(rpc)
ingressStore := getIngressStore(objectStore)
ingressService := NewIngressService(ingressConfig, nodeID, messageBus, ingressClient, ingressRPCClient, ingressStore, roomService, telemetryService)
ioInfoService, err := NewIOInfoService(nodeID, messageBus, egressStore, ingressStore, telemetryService, rpcClient, ingressRPCClient)
if err != nil {
return nil, err
}
rtcService := NewRTCService(conf, roomAllocator, objectStore, router, currentNode, telemetryService)
clientConfigurationManager := createClientConfiguration()
timedVersionGenerator := utils.NewDefaultTimedVersionGenerator()
@@ -92,7 +96,7 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live
if err != nil {
return nil, err
}
livekitServer, err := NewLivekitServer(conf, roomService, egressService, ingressService, rtcService, keyProvider, router, roomManager, server, currentNode)
livekitServer, err := NewLivekitServer(conf, roomService, egressService, ingressService, ioInfoService, rtcService, keyProvider, router, roomManager, server, currentNode)
if err != nil {
return nil, err
}