diff --git a/go.mod b/go.mod index ff1b11db2..df314512c 100644 --- a/go.mod +++ b/go.mod @@ -18,7 +18,7 @@ require ( github.com/jxskiss/base62 v1.1.0 github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 github.com/livekit/mediatransportutil v0.0.0-20230130133657-96cfb115473a - github.com/livekit/protocol v1.5.0 + github.com/livekit/protocol v1.5.1-0.20230314035739-6d1cd857eb3b github.com/livekit/psrpc v0.2.10-0.20230310095745-5cd63568998d github.com/mackerelio/go-osstat v0.2.3 github.com/magefile/mage v1.14.0 diff --git a/go.sum b/go.sum index eb6f39ff8..4e23b6992 100644 --- a/go.sum +++ b/go.sum @@ -233,8 +233,8 @@ github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 h1:jm09419p0lqTkD github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1/go.mod h1:Rs3MhFwutWhGwmY1VQsygw28z5bWcnEYmS1OG9OxjOQ= github.com/livekit/mediatransportutil v0.0.0-20230130133657-96cfb115473a h1:5UkGQpskXp7HcBmyrCwWtO7ygDWbqtjN09Yva4l/nyE= github.com/livekit/mediatransportutil v0.0.0-20230130133657-96cfb115473a/go.mod h1:1Dlx20JPoIKGP45eo+yuj0HjeE25zmyeX/EWHiPCjFw= -github.com/livekit/protocol v1.5.0 h1:jFGSkSEv0PTjUlrW/WnmERejwxyHOSE9If4VU33PYgk= -github.com/livekit/protocol v1.5.0/go.mod h1:hkK/G0wwFiLUGp9F5kxeQxq2CQuIzkmfBwKhTsc71us= +github.com/livekit/protocol v1.5.1-0.20230314035739-6d1cd857eb3b h1:rDm6mdo22cGAFNLwgNhyGwkjxR2civBlspf5//7LIeQ= +github.com/livekit/protocol v1.5.1-0.20230314035739-6d1cd857eb3b/go.mod h1:hkK/G0wwFiLUGp9F5kxeQxq2CQuIzkmfBwKhTsc71us= github.com/livekit/psrpc v0.2.10-0.20230310095745-5cd63568998d h1:3wfbd8zi7zGQCR+xfG3r2k9m2RwXUiIzR0SN4BHewwU= github.com/livekit/psrpc v0.2.10-0.20230310095745-5cd63568998d/go.mod h1:K0j8f1PgLShR7Lx80KbmwFkDH2BvOnycXGV0OSRURKc= github.com/mackerelio/go-osstat v0.2.3 h1:jAMXD5erlDE39kdX2CU7YwCGRcxIO33u/p8+Fhe5dJw= diff --git a/pkg/service/egress.go b/pkg/service/egress.go index abafa7162..1b2c97ec3 100644 --- a/pkg/service/egress.go +++ b/pkg/service/egress.go @@ -300,10 +300,13 @@ func (s *EgressService) ListEgress(ctx context.Context, req *livekit.ListEgressR if err != nil { return nil, err } - items = []*livekit.EgressInfo{info} + + if !req.Active || int32(info.Status) < int32(livekit.EgressStatus_EGRESS_COMPLETE) { + items = []*livekit.EgressInfo{info} + } } else { var err error - items, err = s.es.ListEgress(ctx, livekit.RoomName(req.RoomName)) + items, err = s.es.ListEgress(ctx, livekit.RoomName(req.RoomName), req.Active) if err != nil { return nil, err } diff --git a/pkg/service/interfaces.go b/pkg/service/interfaces.go index 607838f4b..e3b1bcef0 100644 --- a/pkg/service/interfaces.go +++ b/pkg/service/interfaces.go @@ -42,7 +42,7 @@ type ServiceStore interface { type EgressStore interface { StoreEgress(ctx context.Context, info *livekit.EgressInfo) error LoadEgress(ctx context.Context, egressID string) (*livekit.EgressInfo, error) - ListEgress(ctx context.Context, roomName livekit.RoomName) ([]*livekit.EgressInfo, error) + ListEgress(ctx context.Context, roomName livekit.RoomName, active bool) ([]*livekit.EgressInfo, error) UpdateEgress(ctx context.Context, info *livekit.EgressInfo) error } diff --git a/pkg/service/redisstore.go b/pkg/service/redisstore.go index def52e203..2558bc4a4 100644 --- a/pkg/service/redisstore.go +++ b/pkg/service/redisstore.go @@ -350,7 +350,7 @@ func (s *RedisStore) LoadEgress(_ context.Context, egressID string) (*livekit.Eg } } -func (s *RedisStore) ListEgress(_ context.Context, roomName livekit.RoomName) ([]*livekit.EgressInfo, error) { +func (s *RedisStore) ListEgress(_ context.Context, roomName livekit.RoomName, active bool) ([]*livekit.EgressInfo, error) { var infos []*livekit.EgressInfo if roomName == "" { @@ -368,7 +368,11 @@ func (s *RedisStore) ListEgress(_ context.Context, roomName livekit.RoomName) ([ if err != nil { return nil, err } - infos = append(infos, info) + + // if active, filter status starting, active, and ending + if !active || int32(info.Status) < int32(livekit.EgressStatus_EGRESS_COMPLETE) { + infos = append(infos, info) + } } } else { egressIDs, err := s.rc.SMembers(s.ctx, RoomEgressPrefix+string(roomName)).Result() @@ -389,7 +393,11 @@ func (s *RedisStore) ListEgress(_ context.Context, roomName livekit.RoomName) ([ if err != nil { return nil, err } - infos = append(infos, info) + + // if active, filter status starting, active, and ending + if !active || int32(info.Status) < int32(livekit.EgressStatus_EGRESS_COMPLETE) { + infos = append(infos, info) + } } } diff --git a/pkg/service/redisstore_test.go b/pkg/service/redisstore_test.go index 17e8e5021..949480d28 100644 --- a/pkg/service/redisstore_test.go +++ b/pkg/service/redisstore_test.go @@ -189,12 +189,12 @@ func TestEgressStore(t *testing.T) { require.NoError(t, rs.UpdateEgress(ctx, info)) // list - list, err := rs.ListEgress(ctx, "") + list, err := rs.ListEgress(ctx, "", false) require.NoError(t, err) require.Len(t, list, 2) // list by room - list, err = rs.ListEgress(ctx, livekit.RoomName(roomName)) + list, err = rs.ListEgress(ctx, livekit.RoomName(roomName), false) require.NoError(t, err) require.Len(t, list, 1) @@ -207,7 +207,7 @@ func TestEgressStore(t *testing.T) { require.NoError(t, rs.CleanEndedEgress()) // list - list, err = rs.ListEgress(ctx, livekit.RoomName(roomName)) + list, err = rs.ListEgress(ctx, livekit.RoomName(roomName), false) require.NoError(t, err) require.Len(t, list, 0) } diff --git a/pkg/service/servicefakes/fake_egress_store.go b/pkg/service/servicefakes/fake_egress_store.go index 06ee8d9cc..37f13670f 100644 --- a/pkg/service/servicefakes/fake_egress_store.go +++ b/pkg/service/servicefakes/fake_egress_store.go @@ -10,11 +10,12 @@ import ( ) type FakeEgressStore struct { - ListEgressStub func(context.Context, livekit.RoomName) ([]*livekit.EgressInfo, error) + ListEgressStub func(context.Context, livekit.RoomName, bool) ([]*livekit.EgressInfo, error) listEgressMutex sync.RWMutex listEgressArgsForCall []struct { arg1 context.Context arg2 livekit.RoomName + arg3 bool } listEgressReturns struct { result1 []*livekit.EgressInfo @@ -66,19 +67,20 @@ type FakeEgressStore struct { invocationsMutex sync.RWMutex } -func (fake *FakeEgressStore) ListEgress(arg1 context.Context, arg2 livekit.RoomName) ([]*livekit.EgressInfo, error) { +func (fake *FakeEgressStore) ListEgress(arg1 context.Context, arg2 livekit.RoomName, arg3 bool) ([]*livekit.EgressInfo, error) { fake.listEgressMutex.Lock() ret, specificReturn := fake.listEgressReturnsOnCall[len(fake.listEgressArgsForCall)] fake.listEgressArgsForCall = append(fake.listEgressArgsForCall, struct { arg1 context.Context arg2 livekit.RoomName - }{arg1, arg2}) + arg3 bool + }{arg1, arg2, arg3}) stub := fake.ListEgressStub fakeReturns := fake.listEgressReturns - fake.recordInvocation("ListEgress", []interface{}{arg1, arg2}) + fake.recordInvocation("ListEgress", []interface{}{arg1, arg2, arg3}) fake.listEgressMutex.Unlock() if stub != nil { - return stub(arg1, arg2) + return stub(arg1, arg2, arg3) } if specificReturn { return ret.result1, ret.result2 @@ -92,17 +94,17 @@ func (fake *FakeEgressStore) ListEgressCallCount() int { return len(fake.listEgressArgsForCall) } -func (fake *FakeEgressStore) ListEgressCalls(stub func(context.Context, livekit.RoomName) ([]*livekit.EgressInfo, error)) { +func (fake *FakeEgressStore) ListEgressCalls(stub func(context.Context, livekit.RoomName, bool) ([]*livekit.EgressInfo, error)) { fake.listEgressMutex.Lock() defer fake.listEgressMutex.Unlock() fake.ListEgressStub = stub } -func (fake *FakeEgressStore) ListEgressArgsForCall(i int) (context.Context, livekit.RoomName) { +func (fake *FakeEgressStore) ListEgressArgsForCall(i int) (context.Context, livekit.RoomName, bool) { fake.listEgressMutex.RLock() defer fake.listEgressMutex.RUnlock() argsForCall := fake.listEgressArgsForCall[i] - return argsForCall.arg1, argsForCall.arg2 + return argsForCall.arg1, argsForCall.arg2, argsForCall.arg3 } func (fake *FakeEgressStore) ListEgressReturns(result1 []*livekit.EgressInfo, result2 error) {