From aa92a1c71dd2cb8361c37ef7cdd442594e0b31dc Mon Sep 17 00:00:00 2001 From: Mack Swan Date: Fri, 10 Jul 2026 14:55:39 -0400 Subject: [PATCH] Add stream.received webhook event --- .../components/settings/webhook-manager.tsx | 26 ++++++++-- js/components/locales/en-US/settings.ftl | 1 + .../content/docs/lex-reference/openapi.json | 32 ++++++++++-- .../place-stream-server-createwebhook.md | 8 ++- .../server/place-stream-server-defs.md | 8 ++- .../place-stream-server-listwebhooks.md | 20 ++++--- .../place-stream-server-updatewebhook.md | 8 ++- .../place/stream/server/createWebhook.json | 8 ++- lexicons/place/stream/server/defs.json | 8 ++- .../place/stream/server/listWebhooks.json | 8 ++- .../place/stream/server/updateWebhook.json | 8 ++- pkg/director/stream_session.go | 40 +++++++++----- .../discord/send-stream-received.go | 52 +++++++++++++++++++ pkg/integrations/webhook/manager.go | 10 ++++ pkg/statedb/queue_processor.go | 44 ++++++++++++++++ pkg/statedb/webhook_test.go | 18 +++++++ 16 files changed, 264 insertions(+), 35 deletions(-) create mode 100644 pkg/integrations/discord/send-stream-received.go diff --git a/js/app/components/settings/webhook-manager.tsx b/js/app/components/settings/webhook-manager.tsx index f5a9340fd..5cc497b53 100644 --- a/js/app/components/settings/webhook-manager.tsx +++ b/js/app/components/settings/webhook-manager.tsx @@ -75,9 +75,25 @@ interface WebhookFormData { description: string; } +type WebhookEvent = + | "livestream" + | "chat" + | "follow" + | "mention" + | "stream.received"; + +const VALID_WEBHOOK_EVENTS: WebhookEvent[] = [ + "livestream", + "chat", + "follow", + "mention", + "stream.received", +]; + const EVENT_OPTIONS = [ { value: "livestream", labelKey: "events-livestream" }, { value: "chat", labelKey: "events-chat" }, + { value: "stream.received", labelKey: "events-stream-received" }, ]; function WebhookRow({ @@ -631,13 +647,13 @@ export default function WebhookManager() { const response = await agent.place.stream.server.listWebhooks({ limit: 50, }); - // if not type "livestream" | "chat" | "follow" | "mention"[] just return + // Filter out unknown event types returned by the server. // todo: find a better way to check this if (response.data.webhooks) { for (const webhook of response.data.webhooks) { webhook.events = (webhook.events as string[]).filter((event) => - ["livestream", "chat", "follow", "mention"].includes(event), - ) as ("livestream" | "chat" | "follow" | "mention")[]; + VALID_WEBHOOK_EVENTS.includes(event as WebhookEvent), + ) as WebhookEvent[]; } } setWebhooks((response.data.webhooks as any) || []); @@ -663,7 +679,7 @@ export default function WebhookManager() { await agent.place.stream.server.createWebhook({ name: data.name || undefined, url: data.url, - events: data.events as ("livestream" | "chat" | "follow" | "mention")[], + events: data.events as WebhookEvent[], active: data.active, prefix: data.prefix || undefined, suffix: data.suffix || undefined, @@ -700,7 +716,7 @@ export default function WebhookManager() { id: editingWebhook.id, name: data.name || undefined, url: data.url, - events: data.events as ("livestream" | "chat" | "follow" | "mention")[], + events: data.events as WebhookEvent[], active: data.active, prefix: data.prefix || undefined, suffix: data.suffix || undefined, diff --git a/js/components/locales/en-US/settings.ftl b/js/components/locales/en-US/settings.ftl index bd8b366e4..3711dcdc2 100644 --- a/js/components/locales/en-US/settings.ftl +++ b/js/components/locales/en-US/settings.ftl @@ -145,6 +145,7 @@ activates-on = Activates on: events = Events events-livestream = Livestream Events events-chat = Chat Events +events-stream-received = Stream Received Events untitled-webhook = Untitled Webhook inactive = Inactive active = Active diff --git a/js/docs/src/content/docs/lex-reference/openapi.json b/js/docs/src/content/docs/lex-reference/openapi.json index ec5324d5e..43966d2b6 100644 --- a/js/docs/src/content/docs/lex-reference/openapi.json +++ b/js/docs/src/content/docs/lex-reference/openapi.json @@ -619,7 +619,13 @@ "description": "The types of events this webhook should receive.", "items": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"] + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ] } }, "active": { @@ -976,7 +982,13 @@ "schema": { "type": "string", "description": "Filter webhooks that handle this event type.", - "enum": ["chat", "livestream", "follow", "mention"] + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ] } } ] @@ -1059,7 +1071,13 @@ "description": "The types of events this webhook should receive.", "items": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"] + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ] } }, "active": { @@ -6008,7 +6026,13 @@ "description": "The types of events this webhook should receive.", "items": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"] + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ] } }, "active": { diff --git a/js/docs/src/content/docs/lex-reference/server/place-stream-server-createwebhook.md b/js/docs/src/content/docs/lex-reference/server/place-stream-server-createwebhook.md index 60acae40e..9499268c9 100644 --- a/js/docs/src/content/docs/lex-reference/server/place-stream-server-createwebhook.md +++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-createwebhook.md @@ -80,7 +80,13 @@ Create a new webhook for receiving Streamplace events. "type": "array", "items": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"] + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ] }, "description": "The types of events this webhook should receive." }, diff --git a/js/docs/src/content/docs/lex-reference/server/place-stream-server-defs.md b/js/docs/src/content/docs/lex-reference/server/place-stream-server-defs.md index 3d675e8d1..55aa8671c 100644 --- a/js/docs/src/content/docs/lex-reference/server/place-stream-server-defs.md +++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-defs.md @@ -93,7 +93,13 @@ S3 storage configuration for backups. "type": "array", "items": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"] + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ] }, "description": "The types of events this webhook should receive." }, diff --git a/js/docs/src/content/docs/lex-reference/server/place-stream-server-listwebhooks.md b/js/docs/src/content/docs/lex-reference/server/place-stream-server-listwebhooks.md index e0a28e6b8..1c9cc7a66 100644 --- a/js/docs/src/content/docs/lex-reference/server/place-stream-server-listwebhooks.md +++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-listwebhooks.md @@ -17,12 +17,12 @@ List webhooks for the authenticated user. **Parameters:** -| Name | Type | Req'd | Description | Constraints | -| -------- | --------- | ----- | -------------------------------------------- | ----------------------------------------------- | -| `limit` | `integer` | ❌ | The number of webhooks to return. | Min: 1
Max: 100
Default: `50` | -| `cursor` | `string` | ❌ | An optional cursor for pagination. | | -| `active` | `boolean` | ❌ | Filter webhooks by active status. | | -| `event` | `string` | ❌ | Filter webhooks that handle this event type. | Enum: `chat`, `livestream`, `follow`, `mention` | +| Name | Type | Req'd | Description | Constraints | +| -------- | --------- | ----- | -------------------------------------------- | ------------------------------------------------------------------ | +| `limit` | `integer` | ❌ | The number of webhooks to return. | Min: 1
Max: 100
Default: `50` | +| `cursor` | `string` | ❌ | An optional cursor for pagination. | | +| `active` | `boolean` | ❌ | Filter webhooks by active status. | | +| `event` | `string` | ❌ | Filter webhooks that handle this event type. | Enum: `chat`, `livestream`, `follow`, `mention`, `stream.received` | **Output:** @@ -72,7 +72,13 @@ List webhooks for the authenticated user. }, "event": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"], + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ], "description": "Filter webhooks that handle this event type." } } diff --git a/js/docs/src/content/docs/lex-reference/server/place-stream-server-updatewebhook.md b/js/docs/src/content/docs/lex-reference/server/place-stream-server-updatewebhook.md index 67f883fe5..95aa729b1 100644 --- a/js/docs/src/content/docs/lex-reference/server/place-stream-server-updatewebhook.md +++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-updatewebhook.md @@ -86,7 +86,13 @@ Update an existing webhook configuration. "type": "array", "items": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"] + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ] }, "description": "The types of events this webhook should receive." }, diff --git a/lexicons/place/stream/server/createWebhook.json b/lexicons/place/stream/server/createWebhook.json index 146ed4373..93052142b 100644 --- a/lexicons/place/stream/server/createWebhook.json +++ b/lexicons/place/stream/server/createWebhook.json @@ -21,7 +21,13 @@ "type": "array", "items": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"] + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ] }, "description": "The types of events this webhook should receive." }, diff --git a/lexicons/place/stream/server/defs.json b/lexicons/place/stream/server/defs.json index 364f4d694..5d5f0edc7 100644 --- a/lexicons/place/stream/server/defs.json +++ b/lexicons/place/stream/server/defs.json @@ -20,7 +20,13 @@ "type": "array", "items": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"] + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ] }, "description": "The types of events this webhook should receive." }, diff --git a/lexicons/place/stream/server/listWebhooks.json b/lexicons/place/stream/server/listWebhooks.json index 53fcc5ca4..51461d610 100644 --- a/lexicons/place/stream/server/listWebhooks.json +++ b/lexicons/place/stream/server/listWebhooks.json @@ -25,7 +25,13 @@ }, "event": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"], + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ], "description": "Filter webhooks that handle this event type." } } diff --git a/lexicons/place/stream/server/updateWebhook.json b/lexicons/place/stream/server/updateWebhook.json index 4e1ecc4f7..e5a2a20ea 100644 --- a/lexicons/place/stream/server/updateWebhook.json +++ b/lexicons/place/stream/server/updateWebhook.json @@ -25,7 +25,13 @@ "type": "array", "items": { "type": "string", - "enum": ["chat", "livestream", "follow", "mention"] + "enum": [ + "chat", + "livestream", + "follow", + "mention", + "stream.received" + ] }, "description": "The types of events this webhook should receive." }, diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index 8f39be586..8c4a2982f 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -38,18 +38,19 @@ import ( ) type StreamSession struct { - mm *media.MediaManager - mod model.Model - cli *config.CLI - bus *bus.Bus - op *oatproxy.OATProxy - lp *livepeer.LivepeerSession - repoDID string - segmentChan chan struct{} - lastStatus time.Time - lastStatusCID *string - lastOriginTime time.Time - localDB localdb.LocalDB + mm *media.MediaManager + mod model.Model + cli *config.CLI + bus *bus.Bus + op *oatproxy.OATProxy + lp *livepeer.LivepeerSession + repoDID string + segmentChan chan struct{} + lastStatus time.Time + lastStatusCID *string + lastOriginTime time.Time + localDB localdb.LocalDB + streamReceivedNotified bool // Channels for background workers statusUpdateChan chan struct{} // Signal to update status @@ -242,6 +243,21 @@ func (ss *StreamSession) NewSegment(ctx context.Context, notif *media.NewSegment return fmt.Errorf("could not convert segment to streamplace segment: %w", err) } + // stream.received is enqueued once per StreamSession when the first local media segment is accepted. + if notif.Local && !ss.streamReceivedNotified { + ss.streamReceivedNotified = true + ss.Go(ctx, func() error { + task := &statedb.StreamReceivedTask{ + StreamerDID: spseg.Creator, + } + _, err := ss.statefulDB.EnqueueTask(ctx, statedb.TaskStreamReceived, task, statedb.WithTaskKey(fmt.Sprintf("stream-received::%s::%s", spseg.Creator, notif.Segment.ID))) + if err != nil { + log.Error(ctx, "failed to enqueue stream.received task", "err", err) + } + return nil + }) + } + // Record to S3 for live-to-VOD only while the stream is live (an un-ended // place.stream.livestream record exists -> the segment is published). // Segments that arrive before "go live" or after the livestream ends are diff --git a/pkg/integrations/discord/send-stream-received.go b/pkg/integrations/discord/send-stream-received.go new file mode 100644 index 000000000..d12f45996 --- /dev/null +++ b/pkg/integrations/discord/send-stream-received.go @@ -0,0 +1,52 @@ +package discord + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "strings" + + "stream.place/streamplace/pkg/aqhttp" + "stream.place/streamplace/pkg/integrations/discord/discordtypes" +) + +func SendStreamReceived(ctx context.Context, w *discordtypes.Webhook, streamerDID string) error { + content := fmt.Sprintf("stream.received %s", streamerDID) + for _, rewrite := range w.Rewrite { + content = strings.ReplaceAll(content, rewrite.From, rewrite.To) + } + + payload := discordtypes.Payload{ + Content: fmt.Sprintf("%s%s%s", w.Prefix, content, w.Suffix), + } + + jsonPayload, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("failed to marshal payload: %w", err) + } + + req, err := http.NewRequestWithContext(ctx, "POST", w.URL, bytes.NewReader(jsonPayload)) + if err != nil { + return fmt.Errorf("failed to create request: %w", err) + } + req.Header.Set("Content-Type", "application/json") + + resp, err := aqhttp.Do(ctx, req) + if err != nil { + return fmt.Errorf("failed to send request: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusNoContent { + body, err := io.ReadAll(resp.Body) + if err != nil { + return fmt.Errorf("failed to read response body: %w", err) + } + return fmt.Errorf("failed to send request (http %d): %s", resp.StatusCode, string(body)) + } + + return nil +} diff --git a/pkg/integrations/webhook/manager.go b/pkg/integrations/webhook/manager.go index 9da62d797..281cb5667 100644 --- a/pkg/integrations/webhook/manager.go +++ b/pkg/integrations/webhook/manager.go @@ -44,6 +44,16 @@ func SendLivestreamWebhook(ctx context.Context, webhook *streamplace.ServerDefs_ return discord.SendLivestream(ctx, discordWebhook, pdsURL, lsv, postView, spcp) } +// SendStreamReceivedWebhook sends a stream.received event to a specific webhook. +func SendStreamReceivedWebhook(ctx context.Context, webhook *streamplace.ServerDefs_Webhook, streamerDID string) error { + discordWebhook, err := webhookToDiscordWebhook(webhook) + if err != nil { + return fmt.Errorf("failed to convert webhook: %w", err) + } + + return discord.SendStreamReceived(ctx, discordWebhook, streamerDID) +} + // webhookToDiscordWebhook converts streamplace.ServerDefs_Webhook to discordtypes.Webhook func webhookToDiscordWebhook(webhook *streamplace.ServerDefs_Webhook) (*discordtypes.Webhook, error) { var rewriteRules []*discordtypes.WebhookRewrite diff --git a/pkg/statedb/queue_processor.go b/pkg/statedb/queue_processor.go index 927b2b2cd..4e5dd95a5 100644 --- a/pkg/statedb/queue_processor.go +++ b/pkg/statedb/queue_processor.go @@ -24,6 +24,7 @@ import ( var TaskNotification = "notification" var TaskChat = "chat" +var TaskStreamReceived = "stream_received" var TaskFinalizeLivestream = "finalize_livestream" var TaskFinalizeLivestreamVOD = "finalize_livestream_vod" var TaskVODProcess = "vod_process" @@ -38,6 +39,7 @@ var TaskViewCountAggregate = "view_count_aggregate" var nonVODTaskTypes = []string{ TaskNotification, TaskChat, + TaskStreamReceived, TaskFinalizeLivestream, TaskViewCountAggregate, } @@ -53,6 +55,10 @@ type ChatTask struct { MessageView *streamplace.ChatDefs_MessageView } +type StreamReceivedTask struct { + StreamerDID string `json:"streamerDID"` +} + type FinalizeLivestreamTask struct { LivestreamURI string `json:"livestreamURI"` } @@ -154,6 +160,8 @@ func (state *StatefulDB) processTask(ctx context.Context, task *AppTask) error { return state.processNotificationTask(ctx, task) case TaskChat: return state.processChatMessageTask(ctx, task) + case TaskStreamReceived: + return state.processStreamReceivedTask(ctx, task) case TaskFinalizeLivestream: return state.processFinalizeLivestreamTask(ctx, task) case TaskVODProcess: @@ -483,6 +491,42 @@ func (state *StatefulDB) processNotificationTask(ctx context.Context, task *AppT return nil } +func (state *StatefulDB) processStreamReceivedTask(ctx context.Context, task *AppTask) error { + var streamReceivedTask StreamReceivedTask + if err := json.Unmarshal(task.Payload, &streamReceivedTask); err != nil { + return err + } + + webhooks, err := state.GetActiveWebhooksForUser(streamReceivedTask.StreamerDID, "stream.received") + if err != nil { + return fmt.Errorf("failed to get stream.received webhooks: %w", err) + } + for _, w := range webhooks { + lexiconWebhook, err := w.ToLexicon() + if err != nil { + log.Error(ctx, "failed to convert webhook to lexicon", "err", err, "webhook_id", w.ID) + continue + } + go func(lexiconWebhook *streamplace.ServerDefs_Webhook, wid string) { + err := webhook.SendStreamReceivedWebhook(ctx, lexiconWebhook, streamReceivedTask.StreamerDID) + if err != nil { + log.Error(ctx, "failed to send stream.received webhook", "err", err, "webhook_id", wid) + err = state.IncrementWebhookError(wid) + if err != nil { + log.Error(ctx, "failed to increment webhook error count", "err", err, "webhook_id", wid) + } + } else { + log.Log(ctx, "sent stream.received webhook", "webhook_id", wid) + err = state.ResetWebhookError(wid) + if err != nil { + log.Error(ctx, "failed to reset webhook error count", "err", err, "webhook_id", wid) + } + } + }(lexiconWebhook, w.ID) + } + return nil +} + func (state *StatefulDB) processChatMessageTask(ctx context.Context, task *AppTask) error { var chatTask ChatTask if err := json.Unmarshal(task.Payload, &chatTask); err != nil { diff --git a/pkg/statedb/webhook_test.go b/pkg/statedb/webhook_test.go index 8ad6a87c1..bd4ee3f6b 100644 --- a/pkg/statedb/webhook_test.go +++ b/pkg/statedb/webhook_test.go @@ -25,6 +25,24 @@ func TestWebhookCreation(t *testing.T) { }) } +func TestWebhookStreamReceivedEventLookup(t *testing.T) { + WithAllDatabases(t, func(state *StatefulDB) { + webhook := &Webhook{ + UserDID: "did:web:example.com", + URL: "https://example.com", + Events: []byte(`["stream.received"]`), + Active: true, + } + err := state.CreateWebhook(webhook) + require.NoError(t, err) + + activeWebhooks, err := state.GetActiveWebhooksForUser("did:web:example.com", "stream.received") + require.NoError(t, err) + require.Len(t, activeWebhooks, 1) + require.Equal(t, webhook.URL, activeWebhooks[0].URL) + }) +} + func TestWebhookUpdatePartially(t *testing.T) { WithAllDatabases(t, func(state *StatefulDB) { webhook := &Webhook{ -- 2.51.2