From 8c7926e76e225aef964cfe2a65b74e131d66194e Mon Sep 17 00:00:00 2001 From: David Colburn Date: Mon, 13 Jul 2026 15:53:52 -0400 Subject: [PATCH] Egress v2 observability (#1633) * Egress v2 observability * generated protobuf * Fix golangci-lint typecheck error: replace pkg/errors with fmt.Errorf in egressobs --------- Co-authored-by: github-actions <41898282+github-actions[bot]@users.noreply.github.com> Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> --- observability/egressobs/egress.go | 20 +++--- observability/egressobs/egress_test.go | 98 ++++++++++++++++++++++++++ 2 files changed, 108 insertions(+), 10 deletions(-) diff --git a/observability/egressobs/egress.go b/observability/egressobs/egress.go index d72fd48ff..b804a6099 100644 --- a/observability/egressobs/egress.go +++ b/observability/egressobs/egress.go @@ -32,8 +32,8 @@ type EgressResults struct { func GetSourceType(info *livekit.EgressInfo) SessionSourceType { switch r := info.Request.(type) { - // case *livekit.EgressInfo_Egress: - // return getSourceTypeV2(r.Egress) + case *livekit.EgressInfo_Egress: + return getSourceTypeV2(r.Egress) case *livekit.EgressInfo_Replay: return getSourceTypeV2(r.Replay) default: @@ -63,8 +63,8 @@ func getSourceTypeV2(r egress.EgressRequest) SessionSourceType { func GetRequestType(info *livekit.EgressInfo) EgressRequestType { switch info.Request.(type) { - // case *livekit.EgressInfo_Egress: - // return EgressRequestTypeEgress + case *livekit.EgressInfo_Egress: + return EgressRequestTypeEgress case *livekit.EgressInfo_Replay: return EgressRequestTypeReplay case *livekit.EgressInfo_RoomComposite: @@ -105,12 +105,12 @@ func GetStatus(info *livekit.EgressInfo) EgressStatus { func GetRequest(info *livekit.EgressInfo) (string, error) { switch req := info.Request.(type) { - // case *livekit.EgressInfo_Egress: - // b, err := protojson.Marshal(req.Egress) - // if err != nil { - // return "", errors.Wrap(err, "failed to marshal egress request") - // } - // return string(b), nil + case *livekit.EgressInfo_Egress: + b, err := protojson.Marshal(req.Egress) + if err != nil { + return "", fmt.Errorf("failed serializing Egress request: %w", err) + } + return string(b), nil case *livekit.EgressInfo_Replay: b, err := protojson.Marshal(req.Replay) if err != nil { diff --git a/observability/egressobs/egress_test.go b/observability/egressobs/egress_test.go index d99466f0c..4d205cbbf 100644 --- a/observability/egressobs/egress_test.go +++ b/observability/egressobs/egress_test.go @@ -77,6 +77,24 @@ func TestGetRequestType(t *testing.T) { }, expected: "track", }, + { + name: "Egress", + info: &livekit.EgressInfo{ + Request: &livekit.EgressInfo_Egress{ + Egress: &livekit.StartEgressRequest{}, + }, + }, + expected: "egress", + }, + { + name: "Replay", + info: &livekit.EgressInfo{ + Request: &livekit.EgressInfo_Replay{ + Replay: &livekit.ExportReplayRequest{}, + }, + }, + expected: "replay", + }, { name: "Undefined", info: &livekit.EgressInfo{}, @@ -92,6 +110,65 @@ func TestGetRequestType(t *testing.T) { } } +func TestGetSourceTypeV2(t *testing.T) { + tests := []struct { + name string + info *livekit.EgressInfo + expected string + }{ + { + name: "EgressTemplate", + info: &livekit.EgressInfo{ + Request: &livekit.EgressInfo_Egress{ + Egress: &livekit.StartEgressRequest{ + Source: &livekit.StartEgressRequest_Template{Template: &livekit.TemplateSource{}}, + }, + }, + }, + expected: "template", + }, + { + name: "EgressMedia", + info: &livekit.EgressInfo{ + Request: &livekit.EgressInfo_Egress{ + Egress: &livekit.StartEgressRequest{ + Source: &livekit.StartEgressRequest_Media{Media: &livekit.MediaSource{}}, + }, + }, + }, + expected: "media", + }, + { + name: "EgressWeb", + info: &livekit.EgressInfo{ + Request: &livekit.EgressInfo_Egress{ + Egress: &livekit.StartEgressRequest{ + Source: &livekit.StartEgressRequest_Web{Web: &livekit.WebSource{}}, + }, + }, + }, + expected: "web", + }, + { + name: "ReplayTemplate", + info: &livekit.EgressInfo{ + Request: &livekit.EgressInfo_Replay{ + Replay: &livekit.ExportReplayRequest{ + Source: &livekit.ExportReplayRequest_Template{Template: &livekit.TemplateSource{}}, + }, + }, + }, + expected: "template", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + require.Equal(t, tt.expected, string(GetSourceType(tt.info))) + }) + } +} + func TestGetAudioOnly(t *testing.T) { tests := []struct { name string @@ -198,6 +275,27 @@ func TestGetRequest(t *testing.T) { }, }, }, + { + name: "Egress", + info: &livekit.EgressInfo{ + Request: &livekit.EgressInfo_Egress{ + Egress: &livekit.StartEgressRequest{ + RoomName: "test-room", + Source: &livekit.StartEgressRequest_Template{Template: &livekit.TemplateSource{Layout: "speaker"}}, + }, + }, + }, + }, + { + name: "Replay", + info: &livekit.EgressInfo{ + Request: &livekit.EgressInfo_Replay{ + Replay: &livekit.ExportReplayRequest{ + ReplayId: "test-replay", + }, + }, + }, + }, { name: "Undefined", info: &livekit.EgressInfo{},