diff --git a/webhook/resource_url_notifier.go b/webhook/resource_url_notifier.go index 895f9ed76..f941cefe6 100644 --- a/webhook/resource_url_notifier.go +++ b/webhook/resource_url_notifier.go @@ -321,9 +321,12 @@ func (r *ResourceURLNotifier) Process( } sendStart := time.Now() - err := r.send(event, params) + res, err := r.send(event, params) sendDuration := time.Since(sendStart) fields = append(fields, "sendDuration", sendDuration) + if res != nil { + fields = append(fields, "status", res.Status, "statusCode", res.StatusCode) + } if err != nil { params.Logger.Warnw("failed to send webhook", err, fields...) IncDispatchFailure() @@ -349,10 +352,10 @@ func (r *ResourceURLNotifier) Process( } } -func (r *ResourceURLNotifier) send(event *livekit.WebhookEvent, params *ResourceURLNotifierParams) error { +func (r *ResourceURLNotifier) send(event *livekit.WebhookEvent, params *ResourceURLNotifierParams) (*http.Response, error) { encoded, err := protojson.Marshal(event) if err != nil { - return err + return nil, err } // sign payload sum := sha256.Sum256(encoded) @@ -366,22 +369,22 @@ func (r *ResourceURLNotifier) send(event *livekit.WebhookEvent, params *Resource SetSha256(b64) token, err := at.ToJWT() if err != nil { - return err + return nil, err } req, err := retryablehttp.NewRequest("POST", params.URL, bytes.NewReader(encoded)) if err != nil { // ignore and continue - return err + return nil, err } req.Header.Set(authHeader, token) // use a custom mime type to ensure signature is checked prior to parsing req.Header.Set("content-type", "application/webhook+json") res, err := r.client.Do(req) if err != nil { - return err + return nil, err } _ = res.Body.Close() - return nil + return res, nil } func (r *ResourceURLNotifier) sweeper() { diff --git a/webhook/url_notifier.go b/webhook/url_notifier.go index ee96c376f..4512b815a 100644 --- a/webhook/url_notifier.go +++ b/webhook/url_notifier.go @@ -196,9 +196,12 @@ func (n *URLNotifier) QueueNotify(ctx context.Context, event *livekit.WebhookEve fields = append(fields, "queueDuration", queueDuration) sendStart := time.Now() - err := n.send(event, ¶ms) + res, err := n.send(event, ¶ms) sendDuration := time.Since(sendStart) fields = append(fields, "sendDuration", sendDuration) + if res != nil { + fields = append(fields, "status", res.Status, "statusCode", res.StatusCode) + } if err != nil { params.Logger.Warnw("failed to send webhook", err, fields...) n.dropped.Add(event.NumDropped + 1) @@ -258,12 +261,12 @@ func (n *URLNotifier) Stop(force bool) { } } -func (n *URLNotifier) send(event *livekit.WebhookEvent, params *URLNotifierParams) error { +func (n *URLNotifier) send(event *livekit.WebhookEvent, params *URLNotifierParams) (*http.Response, error) { // set dropped count event.NumDropped = n.dropped.Swap(0) encoded, err := protojson.Marshal(event) if err != nil { - return err + return nil, err } // sign payload sum := sha256.Sum256(encoded) @@ -274,20 +277,20 @@ func (n *URLNotifier) send(event *livekit.WebhookEvent, params *URLNotifierParam SetSha256(b64) token, err := at.ToJWT() if err != nil { - return err + return nil, err } r, err := retryablehttp.NewRequest("POST", params.URL, bytes.NewReader(encoded)) if err != nil { // ignore and continue - return err + return nil, err } r.Header.Set(authHeader, token) // use a custom mime type to ensure signature is checked prior to parsing r.Header.Set("content-type", "application/webhook+json") res, err := n.client.Do(r) if err != nil { - return err + return nil, err } _ = res.Body.Close() - return nil + return res, nil } diff --git a/webhook/webhook_test.go b/webhook/webhook_test.go index bf225bc18..c364b6960 100644 --- a/webhook/webhook_test.go +++ b/webhook/webhook_test.go @@ -187,8 +187,9 @@ func TestURLNotifierLifecycle(t *testing.T) { } defer urlNotifier.Stop(false) - err := urlNotifier.send(&livekit.WebhookEvent{Event: EventRoomStarted}, &urlNotifier.params) + res, err := urlNotifier.send(&livekit.WebhookEvent{Event: EventRoomStarted}, &urlNotifier.params) require.Error(t, err) + require.Nil(t, res) }) t.Run("times out before connection", func(t *testing.T) { @@ -208,8 +209,9 @@ func TestURLNotifierLifecycle(t *testing.T) { defer urlNotifier.Stop(false) startedAt := time.Now() - err = urlNotifier.send(&livekit.WebhookEvent{Event: EventRoomStarted}, &urlNotifier.params) + res, err := urlNotifier.send(&livekit.WebhookEvent{Event: EventRoomStarted}, &urlNotifier.params) require.Error(t, err) + require.Nil(t, res) require.Less(t, time.Since(startedAt).Seconds(), float64(2)) }) } @@ -666,8 +668,9 @@ func TestResourceURLNotifierLifecycle(t *testing.T) { } defer resourceURLNotifier.Stop(false) - err := resourceURLNotifier.send(&livekit.WebhookEvent{Event: EventRoomStarted}, ¶ms) + res, err := resourceURLNotifier.send(&livekit.WebhookEvent{Event: EventRoomStarted}, ¶ms) require.Error(t, err) + require.Nil(t, res) }) t.Run("times out before connection", func(t *testing.T) { @@ -694,8 +697,9 @@ func TestResourceURLNotifierLifecycle(t *testing.T) { defer resourceURLNotifier.Stop(false) startedAt := time.Now() - err = resourceURLNotifier.send(&livekit.WebhookEvent{Event: EventRoomStarted}, ¶ms) + res, err := resourceURLNotifier.send(&livekit.WebhookEvent{Event: EventRoomStarted}, ¶ms) require.Error(t, err) + require.Nil(t, res) require.Less(t, time.Since(startedAt).Seconds(), float64(2)) }) }