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{