-
Notifications
You must be signed in to change notification settings - Fork 62
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
* egress * generated protobuf * remove ingress * update egress package * update optional enums * generated protobuf * updated proto * add sent_at to start request * generated protobuf * default codecs * remove extra codecs for now * generated protobuf * put recording back * add connection options * generated protobuf * put recording back * add sent_at and room_id * generated protobuf * add sent_at and room_id * generated protobuf * update egressInfo * put recordingInfo back in webhooks * egress events * update analytics event * egress status * deprecate recording rpcs * undo some changes * key/secret -> token Co-authored-by: github-actions <41898282+github-actions[bot]@users.noreply.github.com>
- Loading branch information
1 parent
28685be
commit ec6020d
Showing
20 changed files
with
5,142 additions
and
745 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,105 @@ | ||
package egress | ||
|
||
import ( | ||
"context" | ||
"errors" | ||
"time" | ||
|
||
"google.golang.org/protobuf/proto" | ||
|
||
"github.com/livekit/protocol/auth" | ||
"github.com/livekit/protocol/livekit" | ||
"github.com/livekit/protocol/logger" | ||
"github.com/livekit/protocol/utils" | ||
) | ||
|
||
const ( | ||
StartChannel = "EG_START" | ||
ResultsChannel = "EG_RESULTS" | ||
requestChannelPrefix = "REQ_" | ||
responseChannelPrefix = "RES_" | ||
LockDuration = time.Second * 3 | ||
requestTimeout = time.Second * 3 | ||
) | ||
|
||
func SendRequest(ctx context.Context, bus utils.MessageBus, req proto.Message) (*livekit.EgressInfo, error) { | ||
requestID := utils.NewGuid(utils.RPCPrefix) | ||
var channel string | ||
|
||
switch r := req.(type) { | ||
case *livekit.StartEgressRequest: | ||
r.EgressId = utils.NewGuid(utils.EgressPrefix) | ||
r.RequestId = requestID | ||
r.SentAt = time.Now().UnixNano() | ||
channel = StartChannel | ||
case *livekit.EgressRequest: | ||
r.RequestId = requestID | ||
channel = RequestChannel(r.EgressId) | ||
default: | ||
return nil, errors.New("invalid request type") | ||
} | ||
|
||
sub, err := bus.Subscribe(ctx, ResponseChannel(requestID)) | ||
if err != nil { | ||
return nil, err | ||
} | ||
defer func() { | ||
err := sub.Close() | ||
if err != nil { | ||
logger.Errorw("failed to unsubscribe from response channel", err) | ||
} | ||
}() | ||
|
||
err = bus.Publish(ctx, channel, req) | ||
if err != nil { | ||
return nil, err | ||
} | ||
|
||
select { | ||
case raw := <-sub.Channel(): | ||
return unmarshalResponse(sub.Payload(raw)) | ||
case <-time.After(requestTimeout): | ||
return nil, errors.New("no response from egress service") | ||
} | ||
} | ||
|
||
func RequestChannel(egressID string) string { | ||
return requestChannelPrefix + egressID | ||
} | ||
|
||
func ResponseChannel(requestID string) string { | ||
return responseChannelPrefix + requestID | ||
} | ||
|
||
func BuildEgressToken(apiKey, secret, roomName string) (string, error) { | ||
f := false | ||
t := true | ||
grant := &auth.VideoGrant{ | ||
RoomJoin: true, | ||
Room: roomName, | ||
CanSubscribe: &t, | ||
CanPublish: &f, | ||
CanPublishData: &f, | ||
Hidden: true, | ||
Recorder: true, | ||
} | ||
|
||
at := auth.NewAccessToken(apiKey, secret). | ||
AddGrant(grant). | ||
SetIdentity(utils.NewGuid(utils.EgressPrefix)). | ||
SetValidFor(24 * time.Hour) | ||
|
||
return at.ToJWT() | ||
} | ||
|
||
func unmarshalResponse(data []byte) (*livekit.EgressInfo, error) { | ||
res := &livekit.EgressResponse{} | ||
err := proto.Unmarshal(data, res) | ||
if err != nil { | ||
return nil, err | ||
} | ||
if res.Error != "" { | ||
return nil, errors.New(res.Error) | ||
} | ||
return res.Info, nil | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.
Oops, something went wrong.
Oops, something went wrong.