diff --git a/js/app/components/settings/settings.tsx b/js/app/components/settings/settings.tsx index 6b30d5409..2a1d82178 100644 --- a/js/app/components/settings/settings.tsx +++ b/js/app/components/settings/settings.tsx @@ -15,6 +15,7 @@ import { Switch } from "react-native"; import { useAppDispatch, useAppSelector } from "store/hooks"; import { Button, H3, H5, Input, Text, View, XStack } from "tamagui"; import { Updates } from "./updates"; +import WebhookManager from "./webhook-manager"; export function Settings() { const dispatch = useAppDispatch(); @@ -144,6 +145,7 @@ export function Settings() { + )} diff --git a/js/app/components/settings/webhook-manager.tsx b/js/app/components/settings/webhook-manager.tsx new file mode 100644 index 000000000..bb48a4365 --- /dev/null +++ b/js/app/components/settings/webhook-manager.tsx @@ -0,0 +1,762 @@ +import { useNavigation } from "@react-navigation/native"; +import { + Button, + Dialog, + DialogFooter, + Input, + Text, + zero, +} from "@streamplace/components"; +import { ThemeProvider } from "@streamplace/components/src/lib/theme/theme"; +import { usePDSAgent } from "@streamplace/components/src/streamplace-store/xrpc"; +import { Edit3, Plus, RefreshCw, Trash2 } from "@tamagui/lucide-icons"; +import AQLink from "components/aqlink"; +import Loading from "components/loading/loading"; +import { useEffect, useState } from "react"; +import { Alert, Pressable, ScrollView, Switch, View } from "react-native"; +import { timeAgo } from "utils/timeAgo"; + +const { + atoms, + bg, + text, + m, + mt, + mr, + mb, + ml, + mx, + my, + p, + pt, + pr, + pb, + pl, + px, + py, + w, + h, + r, + layout, + borders, + flex, + gap, +} = zero; + +interface Webhook { + id: string; + name?: string; + url: string; + events: string[]; + active: boolean; + prefix?: string; + suffix?: string; + description?: string; + createdAt: string; + updatedAt?: string; + lastTriggered?: string; + errorCount?: number; +} + +interface WebhookFormData { + name: string; + url: string; + events: string[]; + active: boolean; + prefix: string; + suffix: string; + description: string; +} + +const EVENT_OPTIONS = [ + { value: "livestream", label: "Livestream Started" }, + { value: "chat", label: "Chat Messages" }, +]; + +function WebhookRow({ + webhook, + onEdit, + onDelete, + isDeleting, +}: { + webhook: Webhook; + onEdit: (webhook: Webhook) => void; + onDelete: (id: string) => void; + isDeleting: boolean; +}) { + const isDiscord = webhook.url + .toLowerCase() + .startsWith("https://discord.com/api/webhooks"); + + return ( + + {/* Header */} + + + + + {webhook.name || "Untitled Webhook"} + + {isDiscord && ( + + Discord + + )} + + + + onEdit(webhook)} + disabled={isDeleting} + > + + + + onDelete(webhook.id)} + disabled={isDeleting} + > + + + + + + {/* Description */} + {webhook.description && ( + + {webhook.description} + + )} + + {/* URL */} + + URL: + + {webhook.url.length > 50 + ? webhook.url.slice(0, 45) + + "..." + + webhook.url.slice(webhook.url.length - 5) + : webhook.url} + + + + {/* Events */} + + Events: + {webhook.events.map((event, index) => ( + + + {EVENT_OPTIONS.find((opt) => opt.value === event)?.label || event} + + + ))} + + + {/* Status info */} + + + Created {timeAgo(new Date(webhook.createdAt))} + + + {webhook.errorCount !== undefined && webhook.errorCount > 0 && ( + + {webhook.errorCount} errors + + )} + {webhook.lastTriggered && ( + + Last triggered {timeAgo(new Date(webhook.lastTriggered))} + + )} + + + + ); +} + +function WebhookForm({ + webhook, + isVisible, + onClose, + onSubmit, + isLoading, +}: { + webhook?: Webhook; + isVisible: boolean; + onClose: () => void; + onSubmit: (data: WebhookFormData) => void; + isLoading: boolean; +}) { + const [formData, setFormData] = useState({ + name: webhook?.name || "", + url: webhook?.url || "", + events: webhook?.events || ["livestream"], + active: webhook?.active ?? true, + prefix: webhook?.prefix || "", + suffix: webhook?.suffix || "", + description: webhook?.description || "", + }); + + const [errors, setErrors] = useState>({}); + + // Update form data when webhook prop changes (for editing) + useEffect(() => { + if (webhook) { + setFormData({ + name: webhook.name || "", + url: webhook.url || "", + events: webhook.events || ["livestream"], + active: webhook.active ?? true, + prefix: webhook.prefix || "", + suffix: webhook.suffix || "", + description: webhook.description || "", + }); + } else { + // Reset form for new webhook + setFormData({ + name: "", + url: "", + events: ["livestream"], + active: true, + prefix: "", + suffix: "", + description: "", + }); + } + }, [webhook]); + + const validateForm = () => { + const newErrors: Record = {}; + + if (!formData.url.trim()) { + newErrors.url = "URL is required"; + } else if (!formData.url.match(/^https?:\/\/.+/)) { + newErrors.url = "URL must start with http:// or https://"; + } + + if (formData.events.length === 0) { + newErrors.events = "At least one event type must be selected"; + } + + setErrors(newErrors); + return Object.keys(newErrors).length === 0; + }; + + const handleSubmit = () => { + if (validateForm()) { + onSubmit(formData); + } + }; + + const toggleEvent = (eventValue: string) => { + setFormData((prev) => ({ + ...prev, + events: prev.events.includes(eventValue) + ? prev.events.filter((e) => e !== eventValue) + : [...prev.events, eventValue], + })); + }; + + return ( + !open && onClose()} + title={webhook ? "Edit Webhook" : "Create Webhook"} + size="lg" + dismissible={false} + > + + {/* Name */} + + + Name (optional) + + + setFormData((prev) => ({ ...prev, name: text })) + } + placeholder="My Discord Webhook" + /> + + + {/* URL */} + + + Webhook URL * + + + setFormData((prev) => ({ ...prev, url: text })) + } + placeholder="https://discord.com/api/webhooks/..." + multiline + /> + {errors.url && ( + + {errors.url} + + )} + + + {/* Description */} + + + Description (optional) + + + setFormData((prev) => ({ ...prev, description: text })) + } + placeholder="Discord notifications for my stream" + multiline + /> + + + {/* Events */} + + + Events * + + {EVENT_OPTIONS.map((option) => ( + toggleEvent(option.value)} + > + + {formData.events.includes(option.value) && ( + ✓ + )} + + + {option.label} + + + ))} + {errors.events && ( + + {errors.events} + + )} + + + {/* Prefix & Suffix */} + + + + Prefix + + + setFormData((prev) => ({ ...prev, prefix: text })) + } + placeholder="🔴 " + /> + + + + Suffix + + + setFormData((prev) => ({ ...prev, suffix: text })) + } + placeholder=" is now live!" + /> + + + + {/* Active toggle */} + + + Active + + + setFormData((prev) => ({ ...prev, active })) + } + /> + + + + + + Cancel + + + + {isLoading ? "Saving..." : webhook ? "Update" : "Create"} + + + + + ); +} + +export default function WebhookManager() { + const navigation = useNavigation(); + const agent = usePDSAgent(); + + const [webhooks, setWebhooks] = useState(null); + const [loading, setLoading] = useState(true); + const [deletingWebhooks, setDeletingWebhooks] = useState>( + new Set(), + ); + const [editingWebhook, setEditingWebhook] = useState(); + const [showForm, setShowForm] = useState(false); + const [formLoading, setFormLoading] = useState(false); + + const loadWebhooks = async () => { + if (!agent) return; + + // wait like 500ms to show loading state + await new Promise((resolve) => setTimeout(resolve, 500)); + + try { + setLoading(true); + const response = await agent.place.stream.server.listWebhooks({ + limit: 50, + }); + // if not type "livestream" | "chat" | "follow" | "mention"[] just return + // 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"[]; + } + } + setWebhooks((response.data.webhooks as any) || []); + } catch (error) { + console.error("Failed to load webhooks:", error); + Alert.alert("Error", "Failed to load webhooks. Please try again."); + } finally { + setLoading(false); + } + }; + + const createWebhook = async (data: WebhookFormData) => { + if (!agent) return; + + try { + setFormLoading(true); + await agent.place.stream.server.createWebhook({ + name: data.name || undefined, + url: data.url, + events: data.events as "livestream" | "chat" | "follow" | "mention"[], + active: data.active, + prefix: data.prefix || undefined, + suffix: data.suffix || undefined, + description: data.description || undefined, + }); + setShowForm(false); + setEditingWebhook(undefined); + await loadWebhooks(); + } catch (error: any) { + console.error("Failed to create webhook:", error); + Alert.alert( + "Error", + error.message || "Failed to create webhook. Please try again.", + ); + } finally { + setFormLoading(false); + } + }; + + const updateWebhook = async (data: WebhookFormData) => { + if (!agent || !editingWebhook) return; + + try { + setFormLoading(true); + await agent.place.stream.server.updateWebhook({ + id: editingWebhook.id, + name: data.name || undefined, + url: data.url, + events: data.events as "livestream" | "chat" | "follow" | "mention"[], + active: data.active, + prefix: data.prefix || undefined, + suffix: data.suffix || undefined, + description: data.description || undefined, + }); + setShowForm(false); + setEditingWebhook(undefined); + await loadWebhooks(); + } catch (error: any) { + console.error("Failed to update webhook:", error); + Alert.alert( + "Error", + error.message || "Failed to update webhook. Please try again.", + ); + } finally { + setFormLoading(false); + } + }; + + const deleteWebhook = async (id: string) => { + if (!agent) return; + + Alert.alert( + "Delete Webhook", + "Are you sure you want to delete this webhook? This action cannot be undone.", + [ + { text: "Cancel", style: "cancel" }, + { + text: "Delete", + style: "destructive", + onPress: async () => { + try { + setDeletingWebhooks((prev) => new Set(prev).add(id)); + await agent.place.stream.server.deleteWebhook({ id }); + await loadWebhooks(); + } catch (error: any) { + console.error("Failed to delete webhook:", error); + Alert.alert( + "Error", + error.message || "Failed to delete webhook. Please try again.", + ); + } finally { + setDeletingWebhooks((prev) => { + const newSet = new Set(prev); + newSet.delete(id); + return newSet; + }); + } + }, + }, + ], + ); + }; + + const handleEdit = (webhook: Webhook) => { + setEditingWebhook(webhook); + setShowForm(true); + }; + + const handleCreate = () => { + setEditingWebhook(undefined); + setShowForm(true); + }; + + const handleSubmit = (data: WebhookFormData) => { + if (editingWebhook) { + updateWebhook(data); + } else { + createWebhook(data); + } + }; + + useEffect(() => { + navigation.setOptions({ title: "Webhook Manager" }); + if (!agent) return; + loadWebhooks(); + }, [agent]); + + return ( + + + + + {/* Header */} + + + Webhook Integrations + + + Create webhooks to receive notifications when you go live or get + chat messages. + + + + + + + + + + {/* Content */} + {loading ? ( + + ) : webhooks === null ? ( + + Failed to load webhooks + + ) : webhooks.length === 0 ? ( + + + No webhooks yet! + + + Create your first webhook to start receiving notifications + when you go live. + + + + Need to set up streaming first? Visit the Live Dashboard + + + + ) : ( + <> + + + {webhooks.length} webhook{webhooks.length !== 1 && "s"} + + + {webhooks.map((webhook) => ( + + ))} + + )} + + + + { + setShowForm(false); + setEditingWebhook(undefined); + }} + onSubmit={handleSubmit} + isLoading={formLoading} + /> + + + ); +} diff --git a/js/docs/src/content/docs/lex-reference/openapi.json b/js/docs/src/content/docs/lex-reference/openapi.json index 84f908a7a..ec559972a 100644 --- a/js/docs/src/content/docs/lex-reference/openapi.json +++ b/js/docs/src/content/docs/lex-reference/openapi.json @@ -6,6 +6,474 @@ "description": "Autogenerated using lexmd" }, "paths": { + "/xrpc/place.stream.server.createWebhook": { + "post": { + "summary": "Create a new webhook for receiving Streamplace events.", + "operationId": "place.stream.server.createWebhook", + "tags": ["place.stream.server"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "webhook": { + "$ref": "#/components/schemas/place.stream.server.defs_webhook" + } + }, + "required": ["webhook"] + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "InvalidUrl" + }, + { + "const": "DuplicateWebhook" + }, + { + "const": "TooManyWebhooks" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "url": { + "type": "string", + "description": "The webhook URL where events will be sent.", + "format": "uri" + }, + "events": { + "type": "array", + "description": "The types of events this webhook should receive.", + "items": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"] + } + }, + "active": { + "type": "boolean", + "description": "Whether this webhook should be active upon creation.", + "default": true + }, + "prefix": { + "type": "string", + "description": "Text to prepend to webhook messages.", + "maxLength": 100 + }, + "suffix": { + "type": "string", + "description": "Text to append to webhook messages.", + "maxLength": 100 + }, + "rewrite": { + "type": "array", + "description": "Text replacement rules for webhook messages.", + "items": { + "$ref": "#/components/schemas/place.stream.server.defs_rewriteRule" + } + }, + "name": { + "type": "string", + "description": "A user-friendly name for this webhook.", + "maxLength": 100 + }, + "description": { + "type": "string", + "description": "A description of what this webhook is used for.", + "maxLength": 500 + } + }, + "required": ["url", "events"] + } + } + } + } + } + }, + "/xrpc/place.stream.server.deleteWebhook": { + "post": { + "summary": "Delete an existing webhook.", + "operationId": "place.stream.server.deleteWebhook", + "tags": ["place.stream.server"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "success": { + "type": "boolean", + "description": "Whether the webhook was successfully deleted." + } + }, + "required": ["success"] + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "WebhookNotFound" + }, + { + "const": "Unauthorized" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "id": { + "type": "string", + "description": "The ID of the webhook to delete." + } + }, + "required": ["id"] + } + } + } + } + } + }, + "/xrpc/place.stream.server.getWebhook": { + "get": { + "summary": "Get details for a specific webhook.", + "operationId": "place.stream.server.getWebhook", + "tags": ["place.stream.server"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "webhook": { + "$ref": "#/components/schemas/place.stream.server.defs_webhook" + } + }, + "required": ["webhook"] + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "WebhookNotFound" + }, + { + "const": "Unauthorized" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "parameters": [ + { + "name": "id", + "in": "query", + "required": true, + "description": "The ID of the webhook to retrieve.", + "schema": { + "type": "string", + "description": "The ID of the webhook to retrieve." + } + } + ] + } + }, + "/xrpc/place.stream.server.listWebhooks": { + "get": { + "summary": "List webhooks for the authenticated user.", + "operationId": "place.stream.server.listWebhooks", + "tags": ["place.stream.server"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "webhooks": { + "type": "array", + "items": { + "$ref": "#/components/schemas/place.stream.server.defs_webhook" + } + }, + "cursor": { + "type": "string", + "description": "A cursor for pagination, if there are more results." + } + }, + "required": ["webhooks"] + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "InvalidCursor" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "parameters": [ + { + "name": "limit", + "in": "query", + "required": false, + "description": "The number of webhooks to return.", + "schema": { + "type": "integer", + "description": "The number of webhooks to return.", + "default": 50, + "minimum": 1, + "maximum": 100 + } + }, + { + "name": "cursor", + "in": "query", + "required": false, + "description": "An optional cursor for pagination.", + "schema": { + "type": "string", + "description": "An optional cursor for pagination." + } + }, + { + "name": "active", + "in": "query", + "required": false, + "description": "Filter webhooks by active status.", + "schema": { + "type": "boolean", + "description": "Filter webhooks by active status." + } + }, + { + "name": "event", + "in": "query", + "required": false, + "description": "Filter webhooks that handle this event type.", + "schema": { + "type": "string", + "description": "Filter webhooks that handle this event type.", + "enum": ["chat", "livestream", "follow", "mention"] + } + } + ] + } + }, + "/xrpc/place.stream.server.updateWebhook": { + "post": { + "summary": "Update an existing webhook configuration.", + "operationId": "place.stream.server.updateWebhook", + "tags": ["place.stream.server"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "webhook": { + "$ref": "#/components/schemas/place.stream.server.defs_webhook" + } + }, + "required": ["webhook"] + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "WebhookNotFound" + }, + { + "const": "Unauthorized" + }, + { + "const": "InvalidUrl" + }, + { + "const": "DuplicateWebhook" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "id": { + "type": "string", + "description": "The ID of the webhook to update." + }, + "url": { + "type": "string", + "description": "The webhook URL where events will be sent.", + "format": "uri" + }, + "events": { + "type": "array", + "description": "The types of events this webhook should receive.", + "items": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"] + } + }, + "active": { + "type": "boolean", + "description": "Whether this webhook should be active." + }, + "prefix": { + "type": "string", + "description": "Text to prepend to webhook messages.", + "maxLength": 100 + }, + "suffix": { + "type": "string", + "description": "Text to append to webhook messages.", + "maxLength": 100 + }, + "rewrite": { + "type": "array", + "description": "Text replacement rules for webhook messages.", + "items": { + "$ref": "#/components/schemas/place.stream.server.defs_rewriteRule" + } + }, + "name": { + "type": "string", + "description": "A user-friendly name for this webhook.", + "maxLength": 100 + }, + "description": { + "type": "string", + "description": "A description of what this webhook is used for.", + "maxLength": 500 + } + }, + "required": ["id"] + } + } + } + } + } + }, "/xrpc/place.stream.live.getLiveUsers": { "get": { "summary": "Get a list of livestream segments for a user", @@ -1034,6 +1502,96 @@ }, "components": { "schemas": { + "place.stream.server.defs_rewriteRule": { + "type": "object", + "properties": { + "from": { + "type": "string", + "description": "Text to search for and replace.", + "maxLength": 100 + }, + "to": { + "type": "string", + "description": "Text to replace with.", + "maxLength": 100 + } + }, + "required": ["from", "to"] + }, + "place.stream.server.defs_webhook": { + "type": "object", + "description": "A webhook configuration for receiving Streamplace events.", + "properties": { + "id": { + "type": "string", + "description": "Unique identifier for this webhook." + }, + "url": { + "type": "string", + "description": "The webhook URL where events will be sent.", + "format": "uri" + }, + "events": { + "type": "array", + "description": "The types of events this webhook should receive.", + "items": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"] + } + }, + "active": { + "type": "boolean", + "description": "Whether this webhook is currently active." + }, + "prefix": { + "type": "string", + "description": "Text to prepend to webhook messages.", + "maxLength": 100 + }, + "suffix": { + "type": "string", + "description": "Text to append to webhook messages.", + "maxLength": 100 + }, + "rewrite": { + "type": "array", + "description": "Text replacement rules for webhook messages.", + "items": { + "$ref": "#/components/schemas/place.stream.server.defs_rewriteRule" + } + }, + "createdAt": { + "type": "string", + "description": "When this webhook was created.", + "format": "date-time" + }, + "updatedAt": { + "type": "string", + "description": "When this webhook was last updated.", + "format": "date-time" + }, + "name": { + "type": "string", + "description": "A user-friendly name for this webhook.", + "maxLength": 100 + }, + "description": { + "type": "string", + "description": "A description of what this webhook is used for.", + "maxLength": 500 + }, + "lastTriggered": { + "type": "string", + "description": "When this webhook was last triggered.", + "format": "date-time" + }, + "errorCount": { + "type": "integer", + "description": "Number of consecutive errors for this webhook." + } + }, + "required": ["id", "url", "events", "active", "createdAt"] + }, "place.stream.livestream_livestreamView": { "type": "object", "properties": { 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 new file mode 100644 index 000000000..65a4b9be8 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-createwebhook.md @@ -0,0 +1,152 @@ +--- +title: place.stream.server.createWebhook +description: Reference for the place.stream.server.createWebhook lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `procedure` + +Create a new webhook for receiving Streamplace events. + +**Parameters:** _(None defined)_ + +**Input:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ------------- | ------------------------------------------------------------------------------------------------------ | ----- | ---------------------------------------------------- | --------------- | +| `url` | `string` | ✅ | The webhook URL where events will be sent. | Format: `uri` | +| `events` | Array of `string` | ✅ | The types of events this webhook should receive. | | +| `active` | `boolean` | ❌ | Whether this webhook should be active upon creation. | Default: `true` | +| `prefix` | `string` | ❌ | Text to prepend to webhook messages. | Max Length: 100 | +| `suffix` | `string` | ❌ | Text to append to webhook messages. | Max Length: 100 | +| `rewrite` | Array of [`place.stream.server.defs#rewriteRule`](/lex-reference/place-stream-server-defs#rewriterule) | ❌ | Text replacement rules for webhook messages. | | +| `name` | `string` | ❌ | A user-friendly name for this webhook. | Max Length: 100 | +| `description` | `string` | ❌ | A description of what this webhook is used for. | Max Length: 500 | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| --------- | ------------------------------------------------------------------------------------- | ----- | ----------- | ----------- | +| `webhook` | [`place.stream.server.defs#webhook`](/lex-reference/place-stream-server-defs#webhook) | ✅ | | | + +**Possible Errors:** + +- `InvalidUrl`: The provided webhook URL is invalid or unreachable. +- `DuplicateWebhook`: A webhook with this URL already exists for this user. +- `TooManyWebhooks`: The user has reached their maximum number of webhooks. + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.server.createWebhook", + "defs": { + "main": { + "type": "procedure", + "description": "Create a new webhook for receiving Streamplace events.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["url", "events"], + "properties": { + "url": { + "type": "string", + "format": "uri", + "description": "The webhook URL where events will be sent." + }, + "events": { + "type": "array", + "items": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"] + }, + "description": "The types of events this webhook should receive." + }, + "active": { + "type": "boolean", + "default": true, + "description": "Whether this webhook should be active upon creation." + }, + "prefix": { + "type": "string", + "maxLength": 100, + "description": "Text to prepend to webhook messages." + }, + "suffix": { + "type": "string", + "maxLength": 100, + "description": "Text to append to webhook messages." + }, + "rewrite": { + "type": "array", + "items": { + "type": "ref", + "ref": "place.stream.server.defs#rewriteRule" + }, + "description": "Text replacement rules for webhook messages." + }, + "name": { + "type": "string", + "maxLength": 100, + "description": "A user-friendly name for this webhook." + }, + "description": { + "type": "string", + "maxLength": 500, + "description": "A description of what this webhook is used for." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["webhook"], + "properties": { + "webhook": { + "type": "ref", + "ref": "place.stream.server.defs#webhook" + } + } + } + }, + "errors": [ + { + "name": "InvalidUrl", + "description": "The provided webhook URL is invalid or unreachable." + }, + { + "name": "DuplicateWebhook", + "description": "A webhook with this URL already exists for this user." + }, + { + "name": "TooManyWebhooks", + "description": "The user has reached their maximum number of webhooks." + } + ] + } + } +} +``` 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 new file mode 100644 index 000000000..1c39f4c33 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-defs.md @@ -0,0 +1,153 @@ +--- +title: place.stream.server.defs +description: Reference for the place.stream.server.defs lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `webhook` + +**Type:** `object` + +A webhook configuration for receiving Streamplace events. + +**Properties:** + +| Name | Type | Req'd | Description | Constraints | +| --------------- | --------------------------------------- | ----- | ------------------------------------------------ | ------------------ | +| `id` | `string` | ✅ | Unique identifier for this webhook. | | +| `url` | `string` | ✅ | The webhook URL where events will be sent. | Format: `uri` | +| `events` | Array of `string` | ✅ | The types of events this webhook should receive. | | +| `active` | `boolean` | ✅ | Whether this webhook is currently active. | | +| `prefix` | `string` | ❌ | Text to prepend to webhook messages. | Max Length: 100 | +| `suffix` | `string` | ❌ | Text to append to webhook messages. | Max Length: 100 | +| `rewrite` | Array of [`#rewriteRule`](#rewriterule) | ❌ | Text replacement rules for webhook messages. | | +| `createdAt` | `string` | ✅ | When this webhook was created. | Format: `datetime` | +| `updatedAt` | `string` | ❌ | When this webhook was last updated. | Format: `datetime` | +| `name` | `string` | ❌ | A user-friendly name for this webhook. | Max Length: 100 | +| `description` | `string` | ❌ | A description of what this webhook is used for. | Max Length: 500 | +| `lastTriggered` | `string` | ❌ | When this webhook was last triggered. | Format: `datetime` | +| `errorCount` | `integer` | ❌ | Number of consecutive errors for this webhook. | | + +--- + + + +### `rewriteRule` + +**Type:** `object` + +**Properties:** + +| Name | Type | Req'd | Description | Constraints | +| ------ | -------- | ----- | ------------------------------- | --------------- | +| `from` | `string` | ✅ | Text to search for and replace. | Max Length: 100 | +| `to` | `string` | ✅ | Text to replace with. | Max Length: 100 | + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.server.defs", + "defs": { + "webhook": { + "type": "object", + "description": "A webhook configuration for receiving Streamplace events.", + "required": ["id", "url", "events", "active", "createdAt"], + "properties": { + "id": { + "type": "string", + "description": "Unique identifier for this webhook." + }, + "url": { + "type": "string", + "format": "uri", + "description": "The webhook URL where events will be sent." + }, + "events": { + "type": "array", + "items": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"] + }, + "description": "The types of events this webhook should receive." + }, + "active": { + "type": "boolean", + "description": "Whether this webhook is currently active." + }, + "prefix": { + "type": "string", + "maxLength": 100, + "description": "Text to prepend to webhook messages." + }, + "suffix": { + "type": "string", + "maxLength": 100, + "description": "Text to append to webhook messages." + }, + "rewrite": { + "type": "array", + "items": { + "type": "ref", + "ref": "#rewriteRule" + }, + "description": "Text replacement rules for webhook messages." + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "When this webhook was created." + }, + "updatedAt": { + "type": "string", + "format": "datetime", + "description": "When this webhook was last updated." + }, + "name": { + "type": "string", + "maxLength": 100, + "description": "A user-friendly name for this webhook." + }, + "description": { + "type": "string", + "maxLength": 500, + "description": "A description of what this webhook is used for." + }, + "lastTriggered": { + "type": "string", + "format": "datetime", + "description": "When this webhook was last triggered." + }, + "errorCount": { + "type": "integer", + "description": "Number of consecutive errors for this webhook." + } + } + }, + "rewriteRule": { + "type": "object", + "required": ["from", "to"], + "properties": { + "from": { + "type": "string", + "maxLength": 100, + "description": "Text to search for and replace." + }, + "to": { + "type": "string", + "maxLength": 100, + "description": "Text to replace with." + } + } + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/server/place-stream-server-deletewebhook.md b/js/docs/src/content/docs/lex-reference/server/place-stream-server-deletewebhook.md new file mode 100644 index 000000000..1036015e9 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-deletewebhook.md @@ -0,0 +1,98 @@ +--- +title: place.stream.server.deleteWebhook +description: Reference for the place.stream.server.deleteWebhook lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `procedure` + +Delete an existing webhook. + +**Parameters:** _(None defined)_ + +**Input:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ---- | -------- | ----- | -------------------------------- | ----------- | +| `id` | `string` | ✅ | The ID of the webhook to delete. | | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| --------- | --------- | ----- | --------------------------------------------- | ----------- | +| `success` | `boolean` | ✅ | Whether the webhook was successfully deleted. | | + +**Possible Errors:** + +- `WebhookNotFound`: The specified webhook was not found. +- `Unauthorized`: The authenticated user does not have access to this webhook. + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.server.deleteWebhook", + "defs": { + "main": { + "type": "procedure", + "description": "Delete an existing webhook.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["id"], + "properties": { + "id": { + "type": "string", + "description": "The ID of the webhook to delete." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["success"], + "properties": { + "success": { + "type": "boolean", + "description": "Whether the webhook was successfully deleted." + } + } + } + }, + "errors": [ + { + "name": "WebhookNotFound", + "description": "The specified webhook was not found." + }, + { + "name": "Unauthorized", + "description": "The authenticated user does not have access to this webhook." + } + ] + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/server/place-stream-server-getwebhook.md b/js/docs/src/content/docs/lex-reference/server/place-stream-server-getwebhook.md new file mode 100644 index 000000000..4f9093578 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-getwebhook.md @@ -0,0 +1,88 @@ +--- +title: place.stream.server.getWebhook +description: Reference for the place.stream.server.getWebhook lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `query` + +Get details for a specific webhook. + +**Parameters:** + +| Name | Type | Req'd | Description | Constraints | +| ---- | -------- | ----- | ---------------------------------- | ----------- | +| `id` | `string` | ✅ | The ID of the webhook to retrieve. | | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| --------- | ------------------------------------------------------------------------------------- | ----- | ----------- | ----------- | +| `webhook` | [`place.stream.server.defs#webhook`](/lex-reference/place-stream-server-defs#webhook) | ✅ | | | + +**Possible Errors:** + +- `WebhookNotFound`: The specified webhook was not found. +- `Unauthorized`: The authenticated user does not have access to this webhook. + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.server.getWebhook", + "defs": { + "main": { + "type": "query", + "description": "Get details for a specific webhook.", + "parameters": { + "type": "params", + "required": ["id"], + "properties": { + "id": { + "type": "string", + "description": "The ID of the webhook to retrieve." + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["webhook"], + "properties": { + "webhook": { + "type": "ref", + "ref": "place.stream.server.defs#webhook" + } + } + } + }, + "errors": [ + { + "name": "WebhookNotFound", + "description": "The specified webhook was not found." + }, + { + "name": "Unauthorized", + "description": "The authenticated user does not have access to this webhook." + } + ] + } + } +} +``` 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 new file mode 100644 index 000000000..e0a28e6b8 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-listwebhooks.md @@ -0,0 +1,109 @@ +--- +title: place.stream.server.listWebhooks +description: Reference for the place.stream.server.listWebhooks lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `query` + +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` | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ---------- | ---------------------------------------------------------------------------------------------- | ----- | --------------------------------------------------- | ----------- | +| `webhooks` | Array of [`place.stream.server.defs#webhook`](/lex-reference/place-stream-server-defs#webhook) | ✅ | | | +| `cursor` | `string` | ❌ | A cursor for pagination, if there are more results. | | + +**Possible Errors:** + +- `InvalidCursor`: The provided cursor is invalid or expired. + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.server.listWebhooks", + "defs": { + "main": { + "type": "query", + "description": "List webhooks for the authenticated user.", + "parameters": { + "type": "params", + "properties": { + "limit": { + "type": "integer", + "minimum": 1, + "maximum": 100, + "default": 50, + "description": "The number of webhooks to return." + }, + "cursor": { + "type": "string", + "description": "An optional cursor for pagination." + }, + "active": { + "type": "boolean", + "description": "Filter webhooks by active status." + }, + "event": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"], + "description": "Filter webhooks that handle this event type." + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["webhooks"], + "properties": { + "webhooks": { + "type": "array", + "items": { + "type": "ref", + "ref": "place.stream.server.defs#webhook" + } + }, + "cursor": { + "type": "string", + "description": "A cursor for pagination, if there are more results." + } + } + } + }, + "errors": [ + { + "name": "InvalidCursor", + "description": "The provided cursor is invalid or expired." + } + ] + } + } +} +``` 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 new file mode 100644 index 000000000..220db4d69 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-updatewebhook.md @@ -0,0 +1,161 @@ +--- +title: place.stream.server.updateWebhook +description: Reference for the place.stream.server.updateWebhook lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `procedure` + +Update an existing webhook configuration. + +**Parameters:** _(None defined)_ + +**Input:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ------------- | ------------------------------------------------------------------------------------------------------ | ----- | ------------------------------------------------ | --------------- | +| `id` | `string` | ✅ | The ID of the webhook to update. | | +| `url` | `string` | ❌ | The webhook URL where events will be sent. | Format: `uri` | +| `events` | Array of `string` | ❌ | The types of events this webhook should receive. | | +| `active` | `boolean` | ❌ | Whether this webhook should be active. | | +| `prefix` | `string` | ❌ | Text to prepend to webhook messages. | Max Length: 100 | +| `suffix` | `string` | ❌ | Text to append to webhook messages. | Max Length: 100 | +| `rewrite` | Array of [`place.stream.server.defs#rewriteRule`](/lex-reference/place-stream-server-defs#rewriterule) | ❌ | Text replacement rules for webhook messages. | | +| `name` | `string` | ❌ | A user-friendly name for this webhook. | Max Length: 100 | +| `description` | `string` | ❌ | A description of what this webhook is used for. | Max Length: 500 | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| --------- | ------------------------------------------------------------------------------------- | ----- | ----------- | ----------- | +| `webhook` | [`place.stream.server.defs#webhook`](/lex-reference/place-stream-server-defs#webhook) | ✅ | | | + +**Possible Errors:** + +- `WebhookNotFound`: The specified webhook was not found. +- `Unauthorized`: The authenticated user does not have access to this webhook. +- `InvalidUrl`: The provided webhook URL is invalid or unreachable. +- `DuplicateWebhook`: A webhook with this URL already exists for this user. + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.server.updateWebhook", + "defs": { + "main": { + "type": "procedure", + "description": "Update an existing webhook configuration.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["id"], + "properties": { + "id": { + "type": "string", + "description": "The ID of the webhook to update." + }, + "url": { + "type": "string", + "format": "uri", + "description": "The webhook URL where events will be sent." + }, + "events": { + "type": "array", + "items": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"] + }, + "description": "The types of events this webhook should receive." + }, + "active": { + "type": "boolean", + "description": "Whether this webhook should be active." + }, + "prefix": { + "type": "string", + "maxLength": 100, + "description": "Text to prepend to webhook messages." + }, + "suffix": { + "type": "string", + "maxLength": 100, + "description": "Text to append to webhook messages." + }, + "rewrite": { + "type": "array", + "items": { + "type": "ref", + "ref": "place.stream.server.defs#rewriteRule" + }, + "description": "Text replacement rules for webhook messages." + }, + "name": { + "type": "string", + "maxLength": 100, + "description": "A user-friendly name for this webhook." + }, + "description": { + "type": "string", + "maxLength": 500, + "description": "A description of what this webhook is used for." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["webhook"], + "properties": { + "webhook": { + "type": "ref", + "ref": "place.stream.server.defs#webhook" + } + } + } + }, + "errors": [ + { + "name": "WebhookNotFound", + "description": "The specified webhook was not found." + }, + { + "name": "Unauthorized", + "description": "The authenticated user does not have access to this webhook." + }, + { + "name": "InvalidUrl", + "description": "The provided webhook URL is invalid or unreachable." + }, + { + "name": "DuplicateWebhook", + "description": "A webhook with this URL already exists for this user." + } + ] + } + } +} +``` diff --git a/lexicons/place/stream/server/createWebhook.json b/lexicons/place/stream/server/createWebhook.json new file mode 100644 index 000000000..4de193389 --- /dev/null +++ b/lexicons/place/stream/server/createWebhook.json @@ -0,0 +1,93 @@ +{ + "lexicon": 1, + "id": "place.stream.server.createWebhook", + "defs": { + "main": { + "type": "procedure", + "description": "Create a new webhook for receiving Streamplace events.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["url", "events"], + "properties": { + "url": { + "type": "string", + "format": "uri", + "description": "The webhook URL where events will be sent." + }, + + "events": { + "type": "array", + "items": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"] + }, + "description": "The types of events this webhook should receive." + }, + "active": { + "type": "boolean", + "default": true, + "description": "Whether this webhook should be active upon creation." + }, + "prefix": { + "type": "string", + "maxLength": 100, + "description": "Text to prepend to webhook messages." + }, + "suffix": { + "type": "string", + "maxLength": 100, + "description": "Text to append to webhook messages." + }, + "rewrite": { + "type": "array", + "items": { + "type": "ref", + "ref": "place.stream.server.defs#rewriteRule" + }, + "description": "Text replacement rules for webhook messages." + }, + "name": { + "type": "string", + "maxLength": 100, + "description": "A user-friendly name for this webhook." + }, + "description": { + "type": "string", + "maxLength": 500, + "description": "A description of what this webhook is used for." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["webhook"], + "properties": { + "webhook": { + "type": "ref", + "ref": "place.stream.server.defs#webhook" + } + } + } + }, + "errors": [ + { + "name": "InvalidUrl", + "description": "The provided webhook URL is invalid or unreachable." + }, + { + "name": "DuplicateWebhook", + "description": "A webhook with this URL already exists for this user." + }, + { + "name": "TooManyWebhooks", + "description": "The user has reached their maximum number of webhooks." + } + ] + } + } +} diff --git a/lexicons/place/stream/server/defs.json b/lexicons/place/stream/server/defs.json new file mode 100644 index 000000000..93977fb61 --- /dev/null +++ b/lexicons/place/stream/server/defs.json @@ -0,0 +1,94 @@ +{ + "lexicon": 1, + "id": "place.stream.server.defs", + "defs": { + "webhook": { + "type": "object", + "description": "A webhook configuration for receiving Streamplace events.", + "required": ["id", "url", "events", "active", "createdAt"], + "properties": { + "id": { + "type": "string", + "description": "Unique identifier for this webhook." + }, + "url": { + "type": "string", + "format": "uri", + "description": "The webhook URL where events will be sent." + }, + "events": { + "type": "array", + "items": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"] + }, + "description": "The types of events this webhook should receive." + }, + "active": { + "type": "boolean", + "description": "Whether this webhook is currently active." + }, + "prefix": { + "type": "string", + "maxLength": 100, + "description": "Text to prepend to webhook messages." + }, + "suffix": { + "type": "string", + "maxLength": 100, + "description": "Text to append to webhook messages." + }, + "rewrite": { + "type": "array", + "items": { "type": "ref", "ref": "#rewriteRule" }, + "description": "Text replacement rules for webhook messages." + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "When this webhook was created." + }, + "updatedAt": { + "type": "string", + "format": "datetime", + "description": "When this webhook was last updated." + }, + "name": { + "type": "string", + "maxLength": 100, + "description": "A user-friendly name for this webhook." + }, + "description": { + "type": "string", + "maxLength": 500, + "description": "A description of what this webhook is used for." + }, + "lastTriggered": { + "type": "string", + "format": "datetime", + "description": "When this webhook was last triggered." + }, + "errorCount": { + "type": "integer", + "description": "Number of consecutive errors for this webhook." + } + } + }, + "rewriteRule": { + "type": "object", + "required": ["from", "to"], + "properties": { + "from": { + "type": "string", + "maxLength": 100, + "description": "Text to search for and replace." + }, + "to": { + "type": "string", + "maxLength": 100, + "description": "Text to replace with." + } + } + } + } +} diff --git a/lexicons/place/stream/server/deleteWebhook.json b/lexicons/place/stream/server/deleteWebhook.json new file mode 100644 index 000000000..fb1868e2c --- /dev/null +++ b/lexicons/place/stream/server/deleteWebhook.json @@ -0,0 +1,46 @@ +{ + "lexicon": 1, + "id": "place.stream.server.deleteWebhook", + "defs": { + "main": { + "type": "procedure", + "description": "Delete an existing webhook.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["id"], + "properties": { + "id": { + "type": "string", + "description": "The ID of the webhook to delete." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["success"], + "properties": { + "success": { + "type": "boolean", + "description": "Whether the webhook was successfully deleted." + } + } + } + }, + "errors": [ + { + "name": "WebhookNotFound", + "description": "The specified webhook was not found." + }, + { + "name": "Unauthorized", + "description": "The authenticated user does not have access to this webhook." + } + ] + } + } +} diff --git a/lexicons/place/stream/server/getWebhook.json b/lexicons/place/stream/server/getWebhook.json new file mode 100644 index 000000000..4fb5fb997 --- /dev/null +++ b/lexicons/place/stream/server/getWebhook.json @@ -0,0 +1,43 @@ +{ + "lexicon": 1, + "id": "place.stream.server.getWebhook", + "defs": { + "main": { + "type": "query", + "description": "Get details for a specific webhook.", + "parameters": { + "type": "params", + "required": ["id"], + "properties": { + "id": { + "type": "string", + "description": "The ID of the webhook to retrieve." + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["webhook"], + "properties": { + "webhook": { + "type": "ref", + "ref": "place.stream.server.defs#webhook" + } + } + } + }, + "errors": [ + { + "name": "WebhookNotFound", + "description": "The specified webhook was not found." + }, + { + "name": "Unauthorized", + "description": "The authenticated user does not have access to this webhook." + } + ] + } + } +} diff --git a/lexicons/place/stream/server/listWebhooks.json b/lexicons/place/stream/server/listWebhooks.json new file mode 100644 index 000000000..53fcc5ca4 --- /dev/null +++ b/lexicons/place/stream/server/listWebhooks.json @@ -0,0 +1,61 @@ +{ + "lexicon": 1, + "id": "place.stream.server.listWebhooks", + "defs": { + "main": { + "type": "query", + "description": "List webhooks for the authenticated user.", + "parameters": { + "type": "params", + "properties": { + "limit": { + "type": "integer", + "minimum": 1, + "maximum": 100, + "default": 50, + "description": "The number of webhooks to return." + }, + "cursor": { + "type": "string", + "description": "An optional cursor for pagination." + }, + "active": { + "type": "boolean", + "description": "Filter webhooks by active status." + }, + "event": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"], + "description": "Filter webhooks that handle this event type." + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["webhooks"], + "properties": { + "webhooks": { + "type": "array", + "items": { + "type": "ref", + "ref": "place.stream.server.defs#webhook" + } + }, + "cursor": { + "type": "string", + "description": "A cursor for pagination, if there are more results." + } + } + } + }, + "errors": [ + { + "name": "InvalidCursor", + "description": "The provided cursor is invalid or expired." + } + ] + } + } +} diff --git a/lexicons/place/stream/server/updateWebhook.json b/lexicons/place/stream/server/updateWebhook.json new file mode 100644 index 000000000..5960a3d6e --- /dev/null +++ b/lexicons/place/stream/server/updateWebhook.json @@ -0,0 +1,100 @@ +{ + "lexicon": 1, + "id": "place.stream.server.updateWebhook", + "defs": { + "main": { + "type": "procedure", + "description": "Update an existing webhook configuration.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["id"], + "properties": { + "id": { + "type": "string", + "description": "The ID of the webhook to update." + }, + "url": { + "type": "string", + "format": "uri", + "description": "The webhook URL where events will be sent." + }, + + "events": { + "type": "array", + "items": { + "type": "string", + "enum": ["chat", "livestream", "follow", "mention"] + }, + "description": "The types of events this webhook should receive." + }, + "active": { + "type": "boolean", + "description": "Whether this webhook should be active." + }, + "prefix": { + "type": "string", + "maxLength": 100, + "description": "Text to prepend to webhook messages." + }, + "suffix": { + "type": "string", + "maxLength": 100, + "description": "Text to append to webhook messages." + }, + "rewrite": { + "type": "array", + "items": { + "type": "ref", + "ref": "place.stream.server.defs#rewriteRule" + }, + "description": "Text replacement rules for webhook messages." + }, + "name": { + "type": "string", + "maxLength": 100, + "description": "A user-friendly name for this webhook." + }, + "description": { + "type": "string", + "maxLength": 500, + "description": "A description of what this webhook is used for." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["webhook"], + "properties": { + "webhook": { + "type": "ref", + "ref": "place.stream.server.defs#webhook" + } + } + } + }, + "errors": [ + { + "name": "WebhookNotFound", + "description": "The specified webhook was not found." + }, + { + "name": "Unauthorized", + "description": "The authenticated user does not have access to this webhook." + }, + { + "name": "InvalidUrl", + "description": "The provided webhook URL is invalid or unreachable." + }, + { + "name": "DuplicateWebhook", + "description": "A webhook with this URL already exists for this user." + } + ] + } + } +} diff --git a/pkg/integrations/discord/avatars.go b/pkg/integrations/discord/avatars.go index 3058649f9..700294810 100644 --- a/pkg/integrations/discord/avatars.go +++ b/pkg/integrations/discord/avatars.go @@ -15,7 +15,7 @@ var avatarCacheMutex = sync.Mutex{} // getAvatarURL gets the avatar URL for a Bluesky from the public appview // pretty ugly. we're going to replace this with indexing bluesky profiles // at some point. -func getAvatarURL(ctx context.Context, did string) (string, error) { +func GetAvatarURL(ctx context.Context, did string) (string, error) { avatarCacheMutex.Lock() defer avatarCacheMutex.Unlock() diff --git a/pkg/integrations/discord/send-chat.go b/pkg/integrations/discord/send-chat.go index d04406483..c0fc7b7cf 100644 --- a/pkg/integrations/discord/send-chat.go +++ b/pkg/integrations/discord/send-chat.go @@ -23,7 +23,7 @@ func SendChat(ctx context.Context, w *discordtypes.Webhook, did string, scm *str return fmt.Errorf("failed to cast chat message to streamplace chat message") } - avatarURL, err := getAvatarURL(ctx, did) + avatarURL, err := GetAvatarURL(ctx, did) if err != nil { log.Warn(ctx, "failed to get avatar URL", "err", err) } diff --git a/pkg/integrations/discord/send-livestream.go b/pkg/integrations/discord/send-livestream.go index dd9398599..414543edc 100644 --- a/pkg/integrations/discord/send-livestream.go +++ b/pkg/integrations/discord/send-livestream.go @@ -37,7 +37,7 @@ func SendLivestream(ctx context.Context, w *discordtypes.Webhook, pdsURL string, Content: fmt.Sprintf("%s%s%s", w.Prefix, content, w.Suffix), } - avatarURL, err := getAvatarURL(ctx, lsv.Author.Did) + avatarURL, err := GetAvatarURL(ctx, lsv.Author.Did) if err != nil { log.Warn(ctx, "failed to get avatar URL", "err", err) } diff --git a/pkg/integrations/webhook/manager.go b/pkg/integrations/webhook/manager.go new file mode 100644 index 000000000..b469ceeff --- /dev/null +++ b/pkg/integrations/webhook/manager.go @@ -0,0 +1,84 @@ +package webhook + +import ( + "context" + "encoding/json" + "fmt" + + "github.com/bluesky-social/indigo/api/bsky" + "gorm.io/datatypes" + "stream.place/streamplace/pkg/integrations/discord" + "stream.place/streamplace/pkg/integrations/discord/discordtypes" + "stream.place/streamplace/pkg/streamplace" +) + +// WebhookData represents the essential webhook information +type WebhookData struct { + ID uint + URL string + Events datatypes.JSON + Active bool + Prefix string + Suffix string + Rewrite datatypes.JSON + Name string +} + +// Manager handles webhook sending +type Manager struct{} + +func NewManager() *Manager { + return &Manager{} +} + +// SendChatWebhook sends chat message to a specific webhook +func (m *Manager) SendChatWebhook(ctx context.Context, webhook WebhookData, authorDID string, scm *streamplace.ChatDefs_MessageView) error { + discordWebhook, err := m.webhookDataToDiscordWebhook(webhook) + if err != nil { + return fmt.Errorf("failed to convert webhook data: %w", err) + } + + return discord.SendChat(ctx, discordWebhook, authorDID, scm) +} + +// SendLivestreamWebhook sends livestream notification to a specific webhook +func (m *Manager) SendLivestreamWebhook(ctx context.Context, webhook WebhookData, pdsURL string, lsv *streamplace.Livestream_LivestreamView, postView *bsky.FeedDefs_PostView, spcp *streamplace.ChatProfile) error { + discordWebhook, err := m.webhookDataToDiscordWebhook(webhook) + if err != nil { + return fmt.Errorf("failed to convert webhook data: %w", err) + } + + return discord.SendLivestream(ctx, discordWebhook, pdsURL, lsv, postView, spcp) +} + +// webhookDataToDiscordWebhook converts WebhookData to discordtypes.Webhook +func (m *Manager) webhookDataToDiscordWebhook(webhook WebhookData) (*discordtypes.Webhook, error) { + var rewriteRules []*discordtypes.WebhookRewrite + if len(webhook.Rewrite) > 0 { + err := json.Unmarshal(webhook.Rewrite, &rewriteRules) + if err != nil { + return nil, fmt.Errorf("failed to unmarshal rewrite rules: %w", err) + } + } + + return &discordtypes.Webhook{ + URL: webhook.URL, + Prefix: webhook.Prefix, + Suffix: webhook.Suffix, + Rewrite: rewriteRules, + }, nil +} + +// WebhookToWebhookData converts a statedb.Webhook to WebhookData to avoid import cycle +func WebhookToWebhookData(id uint, url string, events datatypes.JSON, active bool, prefix, suffix string, rewrite datatypes.JSON, name string) WebhookData { + return WebhookData{ + ID: id, + URL: url, + Events: events, + Active: active, + Prefix: prefix, + Suffix: suffix, + Rewrite: rewrite, + Name: name, + } +} diff --git a/pkg/spxrpc/stubs.go b/pkg/spxrpc/stubs.go index c18676d6e..2d1d8ecc4 100644 --- a/pkg/spxrpc/stubs.go +++ b/pkg/spxrpc/stubs.go @@ -247,6 +247,11 @@ func (s *Server) RegisterHandlersPlaceStream(e *echo.Echo) error { e.GET("/xrpc/place.stream.live.getLiveUsers", s.HandlePlaceStreamLiveGetLiveUsers) e.GET("/xrpc/place.stream.live.getProfileCard", s.HandlePlaceStreamLiveGetProfileCard) e.GET("/xrpc/place.stream.live.getSegments", s.HandlePlaceStreamLiveGetSegments) + e.POST("/xrpc/place.stream.server.createWebhook", s.HandlePlaceStreamServerCreateWebhook) + e.POST("/xrpc/place.stream.server.deleteWebhook", s.HandlePlaceStreamServerDeleteWebhook) + e.GET("/xrpc/place.stream.server.getWebhook", s.HandlePlaceStreamServerGetWebhook) + e.GET("/xrpc/place.stream.server.listWebhooks", s.HandlePlaceStreamServerListWebhooks) + e.POST("/xrpc/place.stream.server.updateWebhook", s.HandlePlaceStreamServerUpdateWebhook) return nil } @@ -330,6 +335,109 @@ func (s *Server) HandlePlaceStreamLiveGetSegments(c echo.Context) error { return c.JSON(200, out) } +func (s *Server) HandlePlaceStreamServerCreateWebhook(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamServerCreateWebhook") + defer span.End() + + var body placestreamtypes.ServerCreateWebhook_Input + if err := c.Bind(&body); err != nil { + return err + } + var out *placestreamtypes.ServerCreateWebhook_Output + var handleErr error + // func (s *Server) handlePlaceStreamServerCreateWebhook(ctx context.Context,body *placestreamtypes.ServerCreateWebhook_Input) (*placestreamtypes.ServerCreateWebhook_Output, error) + out, handleErr = s.handlePlaceStreamServerCreateWebhook(ctx, &body) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + +func (s *Server) HandlePlaceStreamServerDeleteWebhook(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamServerDeleteWebhook") + defer span.End() + + var body placestreamtypes.ServerDeleteWebhook_Input + if err := c.Bind(&body); err != nil { + return err + } + var out *placestreamtypes.ServerDeleteWebhook_Output + var handleErr error + // func (s *Server) handlePlaceStreamServerDeleteWebhook(ctx context.Context,body *placestreamtypes.ServerDeleteWebhook_Input) (*placestreamtypes.ServerDeleteWebhook_Output, error) + out, handleErr = s.handlePlaceStreamServerDeleteWebhook(ctx, &body) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + +func (s *Server) HandlePlaceStreamServerGetWebhook(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamServerGetWebhook") + defer span.End() + id := c.QueryParam("id") + var out *placestreamtypes.ServerGetWebhook_Output + var handleErr error + // func (s *Server) handlePlaceStreamServerGetWebhook(ctx context.Context,id string) (*placestreamtypes.ServerGetWebhook_Output, error) + out, handleErr = s.handlePlaceStreamServerGetWebhook(ctx, id) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + +func (s *Server) HandlePlaceStreamServerListWebhooks(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamServerListWebhooks") + defer span.End() + + var active *bool + if p := c.QueryParam("active"); p != "" { + active_val, err := strconv.ParseBool(p) + if err != nil { + return err + } + active = &active_val + } + cursor := c.QueryParam("cursor") + event := c.QueryParam("event") + + var limit int + if p := c.QueryParam("limit"); p != "" { + var err error + limit, err = strconv.Atoi(p) + if err != nil { + return err + } + } else { + limit = 50 + } + var out *placestreamtypes.ServerListWebhooks_Output + var handleErr error + // func (s *Server) handlePlaceStreamServerListWebhooks(ctx context.Context,active *bool,cursor string,event string,limit int) (*placestreamtypes.ServerListWebhooks_Output, error) + out, handleErr = s.handlePlaceStreamServerListWebhooks(ctx, active, cursor, event, limit) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + +func (s *Server) HandlePlaceStreamServerUpdateWebhook(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamServerUpdateWebhook") + defer span.End() + + var body placestreamtypes.ServerUpdateWebhook_Input + if err := c.Bind(&body); err != nil { + return err + } + var out *placestreamtypes.ServerUpdateWebhook_Output + var handleErr error + // func (s *Server) handlePlaceStreamServerUpdateWebhook(ctx context.Context,body *placestreamtypes.ServerUpdateWebhook_Input) (*placestreamtypes.ServerUpdateWebhook_Output, error) + out, handleErr = s.handlePlaceStreamServerUpdateWebhook(ctx, &body) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + func (s *Server) RegisterHandlersToolsOzone(e *echo.Echo) error { return nil } diff --git a/pkg/spxrpc/webhook.go b/pkg/spxrpc/webhook.go new file mode 100644 index 000000000..700ea6df3 --- /dev/null +++ b/pkg/spxrpc/webhook.go @@ -0,0 +1,357 @@ +package spxrpc + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/url" + "strconv" + "strings" + "time" + + "github.com/labstack/echo/v4" + "github.com/streamplace/oatproxy/pkg/oatproxy" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/statedb" + placestreamtypes "stream.place/streamplace/pkg/streamplace" +) + +func (s *Server) handlePlaceStreamServerCreateWebhook(ctx context.Context, input *placestreamtypes.ServerCreateWebhook_Input) (*placestreamtypes.ServerCreateWebhook_Output, error) { + // Get authenticated user + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session not found") + } + + // Validate input + if input.Url == "" { + return nil, echo.NewHTTPError(http.StatusBadRequest, "URL is required") + } + if len(input.Events) == 0 { + return nil, echo.NewHTTPError(http.StatusBadRequest, "At least one event type is required") + } + + // Validate URL format + if _, err := url.Parse(input.Url); err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, "Invalid URL format") + } + + // Check for duplicate URL for this user + existing, err := s.statefulDB.ListWebhooks(session.DID, 1, 0, map[string]interface{}{ + "url": input.Url, + }) + if err != nil { + log.Error(ctx, "failed to check for duplicate webhook", "err", err) + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to validate webhook") + } + if len(existing) > 0 { + return nil, echo.NewHTTPError(http.StatusConflict, "A webhook with this URL already exists") + } + + // Convert input to database model + eventsJSON, _ := json.Marshal(input.Events) + rewriteJSON, _ := json.Marshal(input.Rewrite) + + webhook := &statedb.Webhook{ + UserDID: session.DID, + URL: input.Url, + Events: eventsJSON, + Active: input.Active != nil && *input.Active, + Prefix: getStringValue(input.Prefix), + Suffix: getStringValue(input.Suffix), + Rewrite: rewriteJSON, + Name: getStringValue(input.Name), + Description: getStringValue(input.Description), + } + + // Create webhook + err = s.statefulDB.CreateWebhook(webhook) + if err != nil { + log.Error(ctx, "failed to create webhook", "err", err) + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to create webhook") + } + + // Convert to API response + apiWebhook, err := dbWebhookToAPI(webhook) + if err != nil { + log.Error(ctx, "failed to convert webhook to API format", "err", err) + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to format webhook response") + } + + return &placestreamtypes.ServerCreateWebhook_Output{ + Webhook: apiWebhook, + }, nil +} + +func (s *Server) handlePlaceStreamServerListWebhooks(ctx context.Context, active *bool, cursor string, event string, limit int) (*placestreamtypes.ServerListWebhooks_Output, error) { + // Get authenticated user + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session not found") + } + + // Set default limit + if limit <= 0 || limit > 100 { + limit = 50 + } + + // Parse cursor for offset + offset := 0 + if cursor != "" { + var err error + offset, err = strconv.Atoi(cursor) + if err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, "Invalid cursor") + } + } + + // Build filters + filters := make(map[string]interface{}) + if active != nil { + filters["active"] = *active + } + + // Get webhooks + webhooks, err := s.statefulDB.ListWebhooks(session.DID, limit+1, offset, filters) + if err != nil { + log.Error(ctx, "failed to list webhooks", "err", err) + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to list webhooks") + } + + // Filter by event type if specified + if event != "" { + filtered := make([]statedb.Webhook, 0) + for _, w := range webhooks { + var events []string + if err := json.Unmarshal(w.Events, &events); err == nil { + for _, e := range events { + if e == event { + filtered = append(filtered, w) + break + } + } + } + } + webhooks = filtered + } + + // Check if there are more results + var nextCursor *string + if len(webhooks) > limit { + webhooks = webhooks[:limit] + next := strconv.Itoa(offset + limit) + nextCursor = &next + } + + // Convert to API format + apiWebhooks := make([]*placestreamtypes.ServerDefs_Webhook, len(webhooks)) + for i, webhook := range webhooks { + apiWebhook, err := dbWebhookToAPI(&webhook) + if err != nil { + log.Error(ctx, "failed to convert webhook to API format", "err", err) + continue + } + apiWebhooks[i] = apiWebhook + } + + return &placestreamtypes.ServerListWebhooks_Output{ + Webhooks: apiWebhooks, + Cursor: nextCursor, + }, nil +} + +func (s *Server) handlePlaceStreamServerGetWebhook(ctx context.Context, id string) (*placestreamtypes.ServerGetWebhook_Output, error) { + // Get authenticated user + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session not found") + } + + // Parse webhook ID + webhookID, err := strconv.ParseUint(id, 10, 32) + if err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, "Invalid webhook ID") + } + + // Get webhook + webhook, err := s.statefulDB.GetWebhook(uint(webhookID), session.DID) + if err != nil { + if strings.Contains(err.Error(), "record not found") { + return nil, echo.NewHTTPError(http.StatusNotFound, "Webhook not found") + } + log.Error(ctx, "failed to get webhook", "err", err) + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to get webhook") + } + + // Convert to API format + apiWebhook, err := dbWebhookToAPI(webhook) + if err != nil { + log.Error(ctx, "failed to convert webhook to API format", "err", err) + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to format webhook response") + } + + return &placestreamtypes.ServerGetWebhook_Output{ + Webhook: apiWebhook, + }, nil +} + +func (s *Server) handlePlaceStreamServerUpdateWebhook(ctx context.Context, input *placestreamtypes.ServerUpdateWebhook_Input) (*placestreamtypes.ServerUpdateWebhook_Output, error) { + // Get authenticated user + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session not found") + } + + // Parse webhook ID + webhookID, err := strconv.ParseUint(input.Id, 10, 32) + if err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, "Invalid webhook ID") + } + + // Validate URL if provided + if input.Url != nil { + if _, err := url.Parse(*input.Url); err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, "Invalid URL format") + } + } + + // Build updates map + updates := make(map[string]interface{}) + if input.Url != nil { + updates["url"] = *input.Url + } + if input.Events != nil { + eventsJSON, _ := json.Marshal(input.Events) + updates["events"] = eventsJSON + } + if input.Active != nil { + updates["active"] = *input.Active + } + if input.Prefix != nil { + updates["prefix"] = *input.Prefix + } + if input.Suffix != nil { + updates["suffix"] = *input.Suffix + } + if input.Rewrite != nil { + rewriteJSON, _ := json.Marshal(input.Rewrite) + updates["rewrite"] = rewriteJSON + } + if input.Name != nil { + updates["name"] = *input.Name + } + if input.Description != nil { + updates["description"] = *input.Description + } + + if len(updates) == 0 { + return nil, echo.NewHTTPError(http.StatusBadRequest, "No fields to update") + } + + // Update webhook + webhook, err := s.statefulDB.UpdateWebhook(uint(webhookID), session.DID, updates) + if err != nil { + if strings.Contains(err.Error(), "record not found") { + return nil, echo.NewHTTPError(http.StatusNotFound, "Webhook not found") + } + log.Error(ctx, "failed to update webhook", "err", err) + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to update webhook") + } + + // Convert to API format + apiWebhook, err := dbWebhookToAPI(webhook) + if err != nil { + log.Error(ctx, "failed to convert webhook to API format", "err", err) + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to format webhook response") + } + + return &placestreamtypes.ServerUpdateWebhook_Output{ + Webhook: apiWebhook, + }, nil +} + +func (s *Server) handlePlaceStreamServerDeleteWebhook(ctx context.Context, input *placestreamtypes.ServerDeleteWebhook_Input) (*placestreamtypes.ServerDeleteWebhook_Output, error) { + // Get authenticated user + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session not found") + } + + // Parse webhook ID + webhookID, err := strconv.ParseUint(input.Id, 10, 32) + if err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, "Invalid webhook ID") + } + + // Delete webhook + err = s.statefulDB.DeleteWebhook(uint(webhookID), session.DID) + if err != nil { + log.Error(ctx, "failed to delete webhook", "err", err) + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to delete webhook") + } + + return &placestreamtypes.ServerDeleteWebhook_Output{ + Success: true, + }, nil +} + +// Helper functions + +func getStringValue(s *string) string { + if s == nil { + return "" + } + return *s +} + +func dbWebhookToAPI(webhook *statedb.Webhook) (*placestreamtypes.ServerDefs_Webhook, error) { + var events []string + if len(webhook.Events) > 0 { + if err := json.Unmarshal(webhook.Events, &events); err != nil { + return nil, fmt.Errorf("failed to unmarshal events: %w", err) + } + } + + var rewrite []*placestreamtypes.ServerDefs_RewriteRule + if len(webhook.Rewrite) > 0 { + if err := json.Unmarshal(webhook.Rewrite, &rewrite); err != nil { + return nil, fmt.Errorf("failed to unmarshal rewrite rules: %w", err) + } + } + + result := &placestreamtypes.ServerDefs_Webhook{ + Id: fmt.Sprintf("%d", webhook.ID), + Url: webhook.URL, + Events: events, + Active: webhook.Active, + CreatedAt: webhook.CreatedAt.Format(time.RFC3339), + ErrorCount: func() *int64 { v := int64(webhook.ErrorCount); return &v }(), + } + + if webhook.Prefix != "" { + result.Prefix = &webhook.Prefix + } + if webhook.Suffix != "" { + result.Suffix = &webhook.Suffix + } + if len(rewrite) > 0 { + result.Rewrite = rewrite + } + if webhook.Name != "" { + result.Name = &webhook.Name + } + if webhook.Description != "" { + result.Description = &webhook.Description + } + if webhook.LastTriggered != nil { + lastTriggered := webhook.LastTriggered.Format(time.RFC3339) + result.LastTriggered = &lastTriggered + } + if !webhook.UpdatedAt.IsZero() { + updatedAt := webhook.UpdatedAt.Format(time.RFC3339) + result.UpdatedAt = &updatedAt + } + + return result, nil +} diff --git a/pkg/statedb/queue_processor.go b/pkg/statedb/queue_processor.go index 733bf1722..694a3f0b5 100644 --- a/pkg/statedb/queue_processor.go +++ b/pkg/statedb/queue_processor.go @@ -9,7 +9,7 @@ import ( "github.com/bluesky-social/indigo/api/bsky" "gorm.io/gorm" - "stream.place/streamplace/pkg/integrations/discord" + "stream.place/streamplace/pkg/integrations/webhook" "stream.place/streamplace/pkg/log" notificationpkg "stream.place/streamplace/pkg/notifications" "stream.place/streamplace/pkg/streamplace" @@ -113,16 +113,30 @@ func (state *StatefulDB) processNotificationTask(ctx context.Context, task *AppT log.Log(ctx, "no notifier configured, skipping notifications", "user", userDID, "count", len(notifications)) } - for _, webhook := range state.CLI.DiscordWebhooks { - if webhook.DID == userDID && webhook.Type == "livestream" { - go func() { - err := discord.SendLivestream(ctx, webhook, notificationTask.PDSURL, lsv, notificationTask.FeedPost, notificationTask.ChatProfile) + // Send to webhooks using webhook manager + webhooks, err := state.GetActiveWebhooksForUser(userDID, "livestream") + if err != nil { + log.Error(ctx, "failed to get livestream webhooks", "err", err) + } else { + webhookManager := webhook.NewManager() + for _, w := range webhooks { + webhookData := webhook.WebhookToWebhookData(w.ID, w.URL, w.Events, w.Active, w.Prefix, w.Suffix, w.Rewrite, w.Name) + go func(wd webhook.WebhookData, wid uint) { + err := webhookManager.SendLivestreamWebhook(ctx, wd, notificationTask.PDSURL, lsv, notificationTask.FeedPost, notificationTask.ChatProfile) if err != nil { - log.Error(ctx, "failed to send livestream to discord", "err", err) + log.Error(ctx, "failed to send livestream to 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 livestream to discord", "user", userDID, "webhook", webhook.URL) + log.Log(ctx, "sent livestream to 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) + } } - }() + }(webhookData, w.ID) } } return nil @@ -138,18 +152,31 @@ func (state *StatefulDB) processChatMessageTask(ctx context.Context, task *AppTa if !ok { return fmt.Errorf("invalid chat message record") } - userDID := scm.Author.Did - for _, webhook := range state.CLI.DiscordWebhooks { - if webhook.DID == rec.Streamer && webhook.Type == "chat" { - go func() { - err := discord.SendChat(ctx, webhook, scm.Author.Did, scm) + // Send to webhooks using webhook manager + webhooks, err := state.GetActiveWebhooksForUser(rec.Streamer, "chat") + if err != nil { + log.Error(ctx, "failed to get chat webhooks", "err", err) + } else { + webhookManager := webhook.NewManager() + for _, w := range webhooks { + webhookData := webhook.WebhookToWebhookData(w.ID, w.URL, w.Events, w.Active, w.Prefix, w.Suffix, w.Rewrite, w.Name) + go func(wd webhook.WebhookData, wid uint) { + err := webhookManager.SendChatWebhook(ctx, wd, scm.Author.Did, scm) if err != nil { - log.Error(ctx, "failed to send livestream to discord", "err", err) + log.Error(ctx, "failed to send chat to 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 livestream to discord", "user", userDID, "webhook", webhook.URL) + log.Log(ctx, "sent chat to 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) + } } - }() + }(webhookData, w.ID) } } return nil diff --git a/pkg/statedb/statedb.go b/pkg/statedb/statedb.go index 28d089caf..dbce09dd9 100644 --- a/pkg/statedb/statedb.go +++ b/pkg/statedb/statedb.go @@ -50,6 +50,7 @@ var StatefulDBModels = []any{ XrpcStreamEvent{}, AppTask{}, Repo{}, + Webhook{}, } var NoPostgresDatabaseCode = "3D000" diff --git a/pkg/statedb/webhook.go b/pkg/statedb/webhook.go new file mode 100644 index 000000000..8fc7e670d --- /dev/null +++ b/pkg/statedb/webhook.go @@ -0,0 +1,93 @@ +package statedb + +import ( + "time" + + "gorm.io/datatypes" +) + +type Webhook struct { + ID uint `gorm:"column:id;primarykey"` + UserDID string `gorm:"column:user_did;not null;index"` + URL string `gorm:"column:url;not null"` + + Events datatypes.JSON `gorm:"column:events;type:jsonb"` + Active bool `gorm:"column:active;default:true"` + Prefix string `gorm:"column:prefix"` + Suffix string `gorm:"column:suffix"` + Rewrite datatypes.JSON `gorm:"column:rewrite;type:jsonb"` + Name string `gorm:"column:name"` + Description string `gorm:"column:description"` + CreatedAt time.Time `gorm:"column:created_at"` + UpdatedAt time.Time `gorm:"column:updated_at"` + LastTriggered *time.Time `gorm:"column:last_triggered"` + ErrorCount int `gorm:"column:error_count;default:0"` +} + +func (w *Webhook) TableName() string { + return "webhooks" +} + +// CreateWebhook creates a new webhook for a user +func (state *StatefulDB) CreateWebhook(webhook *Webhook) error { + return state.DB.Create(webhook).Error +} + +// GetWebhook retrieves a webhook by ID and user DID +func (state *StatefulDB) GetWebhook(id uint, userDID string) (*Webhook, error) { + var webhook Webhook + err := state.DB.Where("id = ? AND user_did = ?", id, userDID).First(&webhook).Error + return &webhook, err +} + +// ListWebhooks retrieves webhooks for a user with optional filters +func (state *StatefulDB) ListWebhooks(userDID string, limit int, offset int, filters map[string]any) ([]Webhook, error) { + var webhooks []Webhook + query := state.DB.Where("user_did = ?", userDID) + + // Apply filters + for key, value := range filters { + if value != nil { + query = query.Where(key+" = ?", value) + } + } + + err := query.Limit(limit).Offset(offset).Order("created_at DESC").Find(&webhooks).Error + return webhooks, err +} + +// UpdateWebhook updates an existing webhook +func (state *StatefulDB) UpdateWebhook(id uint, userDID string, updates map[string]interface{}) (*Webhook, error) { + updates["updated_at"] = time.Now() + err := state.DB.Model(&Webhook{}).Where("id = ? AND user_did = ?", id, userDID).Updates(updates).Error + if err != nil { + return nil, err + } + return state.GetWebhook(id, userDID) +} + +// DeleteWebhook deletes a webhook by ID and user DID +func (state *StatefulDB) DeleteWebhook(id uint, userDID string) error { + return state.DB.Where("id = ? AND user_did = ?", id, userDID).Delete(&Webhook{}).Error +} + +// GetActiveWebhooksForUser retrieves active webhooks for a user filtered by event type +func (state *StatefulDB) GetActiveWebhooksForUser(userDID string, eventType string) ([]Webhook, error) { + var webhooks []Webhook + err := state.DB.Where("user_did = ? AND active = ? AND JSON_EXTRACT(events, '$') LIKE ?", + userDID, true, "%\""+eventType+"\"%").Find(&webhooks).Error + return webhooks, err +} + +// IncrementWebhookError increments the error count for a webhook +func (state *StatefulDB) IncrementWebhookError(id uint) error { + return state.DB.Model(&Webhook{}).Where("id = ?", id).UpdateColumn("error_count", state.DB.Raw("error_count + 1")).Error +} + +// ResetWebhookError resets the error count for a webhook +func (state *StatefulDB) ResetWebhookError(id uint) error { + return state.DB.Model(&Webhook{}).Where("id = ?", id).Updates(map[string]interface{}{ + "error_count": 0, + "last_triggered": time.Now(), + }).Error +} diff --git a/pkg/streamplace/servercreateWebhook.go b/pkg/streamplace/servercreateWebhook.go new file mode 100644 index 000000000..810e06ab4 --- /dev/null +++ b/pkg/streamplace/servercreateWebhook.go @@ -0,0 +1,46 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package streamplace + +// schema: place.stream.server.createWebhook + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +// ServerCreateWebhook_Input is the input argument to a place.stream.server.createWebhook call. +type ServerCreateWebhook_Input struct { + // active: Whether this webhook should be active upon creation. + Active *bool `json:"active,omitempty" cborgen:"active,omitempty"` + // description: A description of what this webhook is used for. + Description *string `json:"description,omitempty" cborgen:"description,omitempty"` + // events: The types of events this webhook should receive. + Events []string `json:"events" cborgen:"events"` + // name: A user-friendly name for this webhook. + Name *string `json:"name,omitempty" cborgen:"name,omitempty"` + // prefix: Text to prepend to webhook messages. + Prefix *string `json:"prefix,omitempty" cborgen:"prefix,omitempty"` + // rewrite: Text replacement rules for webhook messages. + Rewrite []*ServerDefs_RewriteRule `json:"rewrite,omitempty" cborgen:"rewrite,omitempty"` + // suffix: Text to append to webhook messages. + Suffix *string `json:"suffix,omitempty" cborgen:"suffix,omitempty"` + // url: The webhook URL where events will be sent. + Url string `json:"url" cborgen:"url"` +} + +// ServerCreateWebhook_Output is the output of a place.stream.server.createWebhook call. +type ServerCreateWebhook_Output struct { + Webhook *ServerDefs_Webhook `json:"webhook" cborgen:"webhook"` +} + +// ServerCreateWebhook calls the XRPC method "place.stream.server.createWebhook". +func ServerCreateWebhook(ctx context.Context, c util.LexClient, input *ServerCreateWebhook_Input) (*ServerCreateWebhook_Output, error) { + var out ServerCreateWebhook_Output + if err := c.LexDo(ctx, util.Procedure, "application/json", "place.stream.server.createWebhook", nil, input, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/pkg/streamplace/serverdefs.go b/pkg/streamplace/serverdefs.go new file mode 100644 index 000000000..bd87285bb --- /dev/null +++ b/pkg/streamplace/serverdefs.go @@ -0,0 +1,45 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package streamplace + +// schema: place.stream.server.defs + +// ServerDefs_RewriteRule is a "rewriteRule" in the place.stream.server.defs schema. +type ServerDefs_RewriteRule struct { + // from: Text to search for and replace. + From string `json:"from" cborgen:"from"` + // to: Text to replace with. + To string `json:"to" cborgen:"to"` +} + +// ServerDefs_Webhook is a "webhook" in the place.stream.server.defs schema. +// +// A webhook configuration for receiving Streamplace events. +type ServerDefs_Webhook struct { + // active: Whether this webhook is currently active. + Active bool `json:"active" cborgen:"active"` + // createdAt: When this webhook was created. + CreatedAt string `json:"createdAt" cborgen:"createdAt"` + // description: A description of what this webhook is used for. + Description *string `json:"description,omitempty" cborgen:"description,omitempty"` + // errorCount: Number of consecutive errors for this webhook. + ErrorCount *int64 `json:"errorCount,omitempty" cborgen:"errorCount,omitempty"` + // events: The types of events this webhook should receive. + Events []string `json:"events" cborgen:"events"` + // id: Unique identifier for this webhook. + Id string `json:"id" cborgen:"id"` + // lastTriggered: When this webhook was last triggered. + LastTriggered *string `json:"lastTriggered,omitempty" cborgen:"lastTriggered,omitempty"` + // name: A user-friendly name for this webhook. + Name *string `json:"name,omitempty" cborgen:"name,omitempty"` + // prefix: Text to prepend to webhook messages. + Prefix *string `json:"prefix,omitempty" cborgen:"prefix,omitempty"` + // rewrite: Text replacement rules for webhook messages. + Rewrite []*ServerDefs_RewriteRule `json:"rewrite,omitempty" cborgen:"rewrite,omitempty"` + // suffix: Text to append to webhook messages. + Suffix *string `json:"suffix,omitempty" cborgen:"suffix,omitempty"` + // updatedAt: When this webhook was last updated. + UpdatedAt *string `json:"updatedAt,omitempty" cborgen:"updatedAt,omitempty"` + // url: The webhook URL where events will be sent. + Url string `json:"url" cborgen:"url"` +} diff --git a/pkg/streamplace/serverdeleteWebhook.go b/pkg/streamplace/serverdeleteWebhook.go new file mode 100644 index 000000000..6d3ba9929 --- /dev/null +++ b/pkg/streamplace/serverdeleteWebhook.go @@ -0,0 +1,33 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package streamplace + +// schema: place.stream.server.deleteWebhook + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +// ServerDeleteWebhook_Input is the input argument to a place.stream.server.deleteWebhook call. +type ServerDeleteWebhook_Input struct { + // id: The ID of the webhook to delete. + Id string `json:"id" cborgen:"id"` +} + +// ServerDeleteWebhook_Output is the output of a place.stream.server.deleteWebhook call. +type ServerDeleteWebhook_Output struct { + // success: Whether the webhook was successfully deleted. + Success bool `json:"success" cborgen:"success"` +} + +// ServerDeleteWebhook calls the XRPC method "place.stream.server.deleteWebhook". +func ServerDeleteWebhook(ctx context.Context, c util.LexClient, input *ServerDeleteWebhook_Input) (*ServerDeleteWebhook_Output, error) { + var out ServerDeleteWebhook_Output + if err := c.LexDo(ctx, util.Procedure, "application/json", "place.stream.server.deleteWebhook", nil, input, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/pkg/streamplace/servergetWebhook.go b/pkg/streamplace/servergetWebhook.go new file mode 100644 index 000000000..6ba770e0b --- /dev/null +++ b/pkg/streamplace/servergetWebhook.go @@ -0,0 +1,31 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package streamplace + +// schema: place.stream.server.getWebhook + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +// ServerGetWebhook_Output is the output of a place.stream.server.getWebhook call. +type ServerGetWebhook_Output struct { + Webhook *ServerDefs_Webhook `json:"webhook" cborgen:"webhook"` +} + +// ServerGetWebhook calls the XRPC method "place.stream.server.getWebhook". +// +// id: The ID of the webhook to retrieve. +func ServerGetWebhook(ctx context.Context, c util.LexClient, id string) (*ServerGetWebhook_Output, error) { + var out ServerGetWebhook_Output + + params := map[string]interface{}{} + params["id"] = id + if err := c.LexDo(ctx, util.Query, "", "place.stream.server.getWebhook", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/pkg/streamplace/serverlistWebhooks.go b/pkg/streamplace/serverlistWebhooks.go new file mode 100644 index 000000000..aa945d78d --- /dev/null +++ b/pkg/streamplace/serverlistWebhooks.go @@ -0,0 +1,47 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package streamplace + +// schema: place.stream.server.listWebhooks + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +// ServerListWebhooks_Output is the output of a place.stream.server.listWebhooks call. +type ServerListWebhooks_Output struct { + // cursor: A cursor for pagination, if there are more results. + Cursor *string `json:"cursor,omitempty" cborgen:"cursor,omitempty"` + Webhooks []*ServerDefs_Webhook `json:"webhooks" cborgen:"webhooks"` +} + +// ServerListWebhooks calls the XRPC method "place.stream.server.listWebhooks". +// +// active: Filter webhooks by active status. +// cursor: An optional cursor for pagination. +// event: Filter webhooks that handle this event type. +// limit: The number of webhooks to return. +func ServerListWebhooks(ctx context.Context, c util.LexClient, active bool, cursor string, event string, limit int64) (*ServerListWebhooks_Output, error) { + var out ServerListWebhooks_Output + + params := map[string]interface{}{} + if active { + params["active"] = active + } + if cursor != "" { + params["cursor"] = cursor + } + if event != "" { + params["event"] = event + } + if limit != 0 { + params["limit"] = limit + } + if err := c.LexDo(ctx, util.Query, "", "place.stream.server.listWebhooks", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/pkg/streamplace/serverupdateWebhook.go b/pkg/streamplace/serverupdateWebhook.go new file mode 100644 index 000000000..d0bd6e81c --- /dev/null +++ b/pkg/streamplace/serverupdateWebhook.go @@ -0,0 +1,48 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package streamplace + +// schema: place.stream.server.updateWebhook + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +// ServerUpdateWebhook_Input is the input argument to a place.stream.server.updateWebhook call. +type ServerUpdateWebhook_Input struct { + // active: Whether this webhook should be active. + Active *bool `json:"active,omitempty" cborgen:"active,omitempty"` + // description: A description of what this webhook is used for. + Description *string `json:"description,omitempty" cborgen:"description,omitempty"` + // events: The types of events this webhook should receive. + Events []string `json:"events,omitempty" cborgen:"events,omitempty"` + // id: The ID of the webhook to update. + Id string `json:"id" cborgen:"id"` + // name: A user-friendly name for this webhook. + Name *string `json:"name,omitempty" cborgen:"name,omitempty"` + // prefix: Text to prepend to webhook messages. + Prefix *string `json:"prefix,omitempty" cborgen:"prefix,omitempty"` + // rewrite: Text replacement rules for webhook messages. + Rewrite []*ServerDefs_RewriteRule `json:"rewrite,omitempty" cborgen:"rewrite,omitempty"` + // suffix: Text to append to webhook messages. + Suffix *string `json:"suffix,omitempty" cborgen:"suffix,omitempty"` + // url: The webhook URL where events will be sent. + Url *string `json:"url,omitempty" cborgen:"url,omitempty"` +} + +// ServerUpdateWebhook_Output is the output of a place.stream.server.updateWebhook call. +type ServerUpdateWebhook_Output struct { + Webhook *ServerDefs_Webhook `json:"webhook" cborgen:"webhook"` +} + +// ServerUpdateWebhook calls the XRPC method "place.stream.server.updateWebhook". +func ServerUpdateWebhook(ctx context.Context, c util.LexClient, input *ServerUpdateWebhook_Input) (*ServerUpdateWebhook_Output, error) { + var out ServerUpdateWebhook_Output + if err := c.LexDo(ctx, util.Procedure, "application/json", "place.stream.server.updateWebhook", nil, input, &out); err != nil { + return nil, err + } + + return &out, nil +}