From 552e3758d5535a41218944359320a758e9ba9e9d Mon Sep 17 00:00:00 2001 From: Benjamin Pracht Date: Fri, 16 Jun 2023 10:58:49 -0700 Subject: [PATCH] Add IngressUpdated event (#1775) --- pkg/service/ioinfo.go | 5 +++ pkg/telemetry/events.go | 6 +++ .../telemetryfakes/fake_telemetry_service.go | 41 +++++++++++++++++++ pkg/telemetry/telemetryservice.go | 1 + 4 files changed, 53 insertions(+) diff --git a/pkg/service/ioinfo.go b/pkg/service/ioinfo.go index 41ff4541c..f3bd6cac7 100644 --- a/pkg/service/ioinfo.go +++ b/pkg/service/ioinfo.go @@ -137,6 +137,11 @@ func (s *IOInfoService) UpdateIngressState(ctx context.Context, req *rpc.UpdateI s.telemetry.IngressStarted(ctx, info) logger.Infow("ingress started", "ingressID", req.IngressId) + + case livekit.IngressState_ENDPOINT_BUFFERING: + s.telemetry.IngressUpdated(ctx, info) + + logger.Infow("ingress buffering", "ingressID", req.IngressId) } } diff --git a/pkg/telemetry/events.go b/pkg/telemetry/events.go index 71d8c05b7..b1e731bdc 100644 --- a/pkg/telemetry/events.go +++ b/pkg/telemetry/events.go @@ -456,6 +456,12 @@ func (t *telemetryService) IngressStarted(ctx context.Context, info *livekit.Ing }) } +func (t *telemetryService) IngressUpdated(ctx context.Context, info *livekit.IngressInfo) { + t.enqueue(func() { + t.SendEvent(ctx, newIngressEvent(livekit.AnalyticsEventType_INGRESS_UPDATED, info)) + }) +} + func (t *telemetryService) IngressEnded(ctx context.Context, info *livekit.IngressInfo) { t.enqueue(func() { t.NotifyEvent(ctx, &livekit.WebhookEvent{ diff --git a/pkg/telemetry/telemetryfakes/fake_telemetry_service.go b/pkg/telemetry/telemetryfakes/fake_telemetry_service.go index f85e03527..fd3ae6ff5 100644 --- a/pkg/telemetry/telemetryfakes/fake_telemetry_service.go +++ b/pkg/telemetry/telemetryfakes/fake_telemetry_service.go @@ -56,6 +56,12 @@ type FakeTelemetryService struct { arg1 context.Context arg2 *livekit.IngressInfo } + IngressUpdatedStub func(context.Context, *livekit.IngressInfo) + ingressUpdatedMutex sync.RWMutex + ingressUpdatedArgsForCall []struct { + arg1 context.Context + arg2 *livekit.IngressInfo + } NotifyEventStub func(context.Context, *livekit.WebhookEvent) notifyEventMutex sync.RWMutex notifyEventArgsForCall []struct { @@ -494,6 +500,39 @@ func (fake *FakeTelemetryService) IngressStartedArgsForCall(i int) (context.Cont return argsForCall.arg1, argsForCall.arg2 } +func (fake *FakeTelemetryService) IngressUpdated(arg1 context.Context, arg2 *livekit.IngressInfo) { + fake.ingressUpdatedMutex.Lock() + fake.ingressUpdatedArgsForCall = append(fake.ingressUpdatedArgsForCall, struct { + arg1 context.Context + arg2 *livekit.IngressInfo + }{arg1, arg2}) + stub := fake.IngressUpdatedStub + fake.recordInvocation("IngressUpdated", []interface{}{arg1, arg2}) + fake.ingressUpdatedMutex.Unlock() + if stub != nil { + fake.IngressUpdatedStub(arg1, arg2) + } +} + +func (fake *FakeTelemetryService) IngressUpdatedCallCount() int { + fake.ingressUpdatedMutex.RLock() + defer fake.ingressUpdatedMutex.RUnlock() + return len(fake.ingressUpdatedArgsForCall) +} + +func (fake *FakeTelemetryService) IngressUpdatedCalls(stub func(context.Context, *livekit.IngressInfo)) { + fake.ingressUpdatedMutex.Lock() + defer fake.ingressUpdatedMutex.Unlock() + fake.IngressUpdatedStub = stub +} + +func (fake *FakeTelemetryService) IngressUpdatedArgsForCall(i int) (context.Context, *livekit.IngressInfo) { + fake.ingressUpdatedMutex.RLock() + defer fake.ingressUpdatedMutex.RUnlock() + argsForCall := fake.ingressUpdatedArgsForCall[i] + return argsForCall.arg1, argsForCall.arg2 +} + func (fake *FakeTelemetryService) NotifyEvent(arg1 context.Context, arg2 *livekit.WebhookEvent) { fake.notifyEventMutex.Lock() fake.notifyEventArgsForCall = append(fake.notifyEventArgsForCall, struct { @@ -1318,6 +1357,8 @@ func (fake *FakeTelemetryService) Invocations() map[string][][]interface{} { defer fake.ingressEndedMutex.RUnlock() fake.ingressStartedMutex.RLock() defer fake.ingressStartedMutex.RUnlock() + fake.ingressUpdatedMutex.RLock() + defer fake.ingressUpdatedMutex.RUnlock() fake.notifyEventMutex.RLock() defer fake.notifyEventMutex.RUnlock() fake.participantActiveMutex.RLock() diff --git a/pkg/telemetry/telemetryservice.go b/pkg/telemetry/telemetryservice.go index e304980aa..5b4c4b9c1 100644 --- a/pkg/telemetry/telemetryservice.go +++ b/pkg/telemetry/telemetryservice.go @@ -57,6 +57,7 @@ type TelemetryService interface { IngressCreated(ctx context.Context, info *livekit.IngressInfo) IngressDeleted(ctx context.Context, info *livekit.IngressInfo) IngressStarted(ctx context.Context, info *livekit.IngressInfo) + IngressUpdated(ctx context.Context, info *livekit.IngressInfo) IngressEnded(ctx context.Context, info *livekit.IngressInfo) // helpers