diff --git a/js/app/components/mobile/desktop-ui.tsx b/js/app/components/mobile/desktop-ui.tsx index afb0d6ebc..ed8af6ed6 100644 --- a/js/app/components/mobile/desktop-ui.tsx +++ b/js/app/components/mobile/desktop-ui.tsx @@ -18,13 +18,13 @@ import Animated, { useSharedValue, withTiming, } from "react-native-reanimated"; +import { useSafeAreaInsets } from "react-native-safe-area-context"; import { BottomControlBar, MuteOverlay, TopControlBar, } from "./desktop-ui/index"; import { useResponsiveLayout } from "./useResponsiveLayout"; -import { useSafeAreaInsets } from "react-native-safe-area-context"; const { h, layout, position, w, px, py, r, p } = zero; @@ -220,7 +220,6 @@ export function DesktopUi({ ingest={ingest} isChatOpen={isChatOpen || false} onToggleChat={toggleChat} - safeAreaInsets={safeAreaInsets} embedded={embedded} /> diff --git a/js/app/components/mobile/player.tsx b/js/app/components/mobile/player.tsx index 73e9909c5..6035f7a24 100644 --- a/js/app/components/mobile/player.tsx +++ b/js/app/components/mobile/player.tsx @@ -13,7 +13,6 @@ import { usePlayerDimensions, usePlayerStore, useSegment, - useSegmentDimensions, View, } from "@streamplace/components"; import { gap, h, pt, w } from "@streamplace/components/src/lib/theme/atoms"; @@ -27,7 +26,7 @@ import Reanimated, { useSharedValue, withTiming, } from "react-native-reanimated"; -import { useAppSelector } from "store/hooks"; +import { useUserProfile } from "store/hooks"; import { BottomMetadata } from "./bottom-metadata"; import { DesktopChatPanel } from "./chat"; import { DesktopUi } from "./desktop-ui"; @@ -35,11 +34,8 @@ import { OfflineCounter } from "./offline-counter"; import { MobileUi } from "./ui"; import { useResponsiveLayout } from "./useResponsiveLayout"; -import { - setSidebarHidden, - setSidebarUnhidden, -} from "features/base/sidebarSlice"; -import { useDispatch } from "react-redux"; +import { useSafeAreaInsets } from "react-native-safe-area-context"; +import { useStore } from "store"; import { UserOffline } from "./user-offline"; const SEGMENT_TIMEOUT = 500; // half a sec @@ -225,6 +221,10 @@ export function PlayerInner( showChatSidePanelOnLandscape: props.showChat, }); + const safeAreaInsets = useSafeAreaInsets(); + const setSidebarHidden = useStore((state) => state.setSidebarHidden); + const setSidebarUnhidden = useStore((state) => state.setSidebarUnhidden); + // auto-collapse chat once when going offline const hasCollapsedChat = useRef(false); useEffect(() => { @@ -252,9 +252,6 @@ export function PlayerInner( } }, [props.showUnavailable]); - // for hiding sidebar - const dispatch = useDispatch(); - // content info const { width, height } = usePlayerDimensions(); @@ -291,6 +288,8 @@ export function PlayerInner( const showFullDesktopMode = aspectRatio > 1 && screenWidth > 1200; const isLandscape = aspectRatio > 1; + const isPlayerRatioGreater = aspectRatio >= 16 / 9; + // animated style for offline height transition const animatedHeightStyle = useAnimatedStyle(() => { return { diff --git a/js/app/components/settings/recommendations-manager.tsx b/js/app/components/settings/recommendations-manager.tsx new file mode 100644 index 000000000..6ef2cd86b --- /dev/null +++ b/js/app/components/settings/recommendations-manager.tsx @@ -0,0 +1,567 @@ +import { + Button, + Dialog, + Input, + MenuContainer, + MenuGroup, + MenuInfo, + MenuSeparator, + Text, + zero, +} from "@streamplace/components"; +import { usePDSAgent } from "@streamplace/components/src/streamplace-store/xrpc"; +import Loading from "components/loading/loading"; +import { Plus, RefreshCw, Search, X } from "lucide-react-native"; +import { useCallback, useEffect, useState } from "react"; +import { useTranslation } from "react-i18next"; +import { Alert, Pressable, ScrollView, View } from "react-native"; + +const { text, mt, mb, px, py, w, layout, gap, r, p } = zero; + +interface ActorSearchResult { + did: string; + handle: string; +} + +export default function RecommendationsManager() { + const agent = usePDSAgent(); + const { theme } = zero.useTheme(); + const [streamers, setStreamers] = useState([]); + const [loading, setLoading] = useState(true); + const [saving, setSaving] = useState(false); + const [deleteDialog, setDeleteDialog] = useState<{ + isVisible: boolean; + index: number | null; + }>({ isVisible: false, index: null }); + const [errors, setErrors] = useState>({}); + + // Search state + const [searchQuery, setSearchQuery] = useState(""); + const [searchResults, setSearchResults] = useState([]); + const [searching, setSearching] = useState(false); + const [searchDebounceTimeout, setSearchDebounceTimeout] = + useState(null); + + const { t } = useTranslation("settings"); + + const loadRecommendations = async () => { + if (!agent) return; + + try { + setLoading(true); + const userDID = agent.did; + if (!userDID) { + setStreamers([]); + return; + } + + // Get the record directly from the PDS for editing + const response = await agent.com.atproto.repo.getRecord({ + repo: userDID, + collection: "place.stream.live.recommendations", + rkey: "self", + }); + + const record = response.data.value as any; + setStreamers(record.streamers || []); + } catch (error: any) { + console.error("Failed to load recommendations:", error); + if (error.status !== 404) { + Alert.alert( + "Error", + "Failed to load recommendations. Please try again.", + ); + } + setStreamers([]); + } finally { + setLoading(false); + } + }; + + const saveRecommendations = async (newStreamers: string[]) => { + if (!agent || saving) return; + + try { + if (!agent.did) { + throw new Error("Agent DID is not available"); + } + setSaving(true); + + await agent.place.stream.live.recommendations.create( + { + repo: agent.did, + }, + { + createdAt: new Date().toISOString(), + streamers: newStreamers, + }, + ); + + setStreamers(newStreamers); + } catch (error: any) { + console.error("Failed to save recommendations:", error); + Alert.alert( + "Error", + error.message || "Failed to save recommendations. Please try again.", + ); + // Reload to get back to consistent state + await loadRecommendations(); + } finally { + setSaving(false); + } + }; + + const searchActors = useCallback( + async (query: string) => { + if (!agent || !query.trim()) { + setSearchResults([]); + return; + } + + try { + setSearching(true); + const response = await agent.place.stream.live.searchActorsTypeahead({ + q: query, + limit: 10, + }); + + setSearchResults( + response.data.actors.map((actor: any) => ({ + did: actor.did, + handle: actor.handle, + })), + ); + } catch (error: any) { + console.error("Failed to search actors:", error); + setSearchResults([]); + } finally { + setSearching(false); + } + }, + [agent], + ); + + const handleSearchChange = (query: string) => { + setSearchQuery(query); + + // Clear previous timeout + if (searchDebounceTimeout) { + clearTimeout(searchDebounceTimeout); + } + + // Set new timeout for debounced search + if (query.trim()) { + const timeout = setTimeout(() => { + searchActors(query); + }, 300); + setSearchDebounceTimeout(timeout); + } else { + setSearchResults([]); + } + }; + + const handleSelectActor = async (actor: ActorSearchResult) => { + if (streamers.length >= 8) { + Alert.alert( + "Maximum Reached", + "You can only add up to 8 recommendations.", + ); + return; + } + + if (streamers.includes(actor.did)) { + Alert.alert( + "Already Added", + "This streamer is already in your recommendations.", + ); + return; + } + + const newStreamers = [...streamers, actor.did]; + await saveRecommendations(newStreamers); + + // Clear search + setSearchQuery(""); + setSearchResults([]); + }; + + const validateDID = (did: string, index: number): boolean => { + const trimmed = did.trim(); + if (!trimmed) { + setErrors((prev) => ({ ...prev, [index]: "DID is required" })); + return false; + } + if (!trimmed.startsWith("did:")) { + setErrors((prev) => ({ + ...prev, + [index]: "DID must start with 'did:'", + })); + return false; + } + setErrors((prev) => { + const newErrors = { ...prev }; + delete newErrors[index]; + return newErrors; + }); + return true; + }; + + const handleStreamerChange = async (index: number, value: string) => { + const newStreamers = [...streamers]; + newStreamers[index] = value; + setStreamers(newStreamers); + }; + + const handleStreamerBlur = async (index: number) => { + const value = streamers[index].trim(); + if (!value) { + // Empty field, just remove it + const newStreamers = streamers.filter((_, i) => i !== index); + await saveRecommendations(newStreamers); + return; + } + + if (validateDID(value, index)) { + await saveRecommendations(streamers); + } + }; + + const handleAddRecommendation = () => { + if (streamers.length >= 8) { + Alert.alert( + "Maximum Reached", + "You can only add up to 8 recommendations.", + ); + return; + } + setStreamers([...streamers, ""]); + }; + + const handleDelete = (index: number) => { + setDeleteDialog({ isVisible: true, index }); + }; + + const confirmDelete = async () => { + if (deleteDialog.index === null) return; + + const newStreamers = streamers.filter((_, i) => i !== deleteDialog.index); + await saveRecommendations(newStreamers); + setDeleteDialog({ isVisible: false, index: null }); + }; + + useEffect(() => { + if (!agent) return; + loadRecommendations(); + }, [agent]); + + // Cleanup timeout on unmount + useEffect(() => { + return () => { + if (searchDebounceTimeout) { + clearTimeout(searchDebounceTimeout); + } + }; + }, [searchDebounceTimeout]); + + if (!agent) { + return ; + } + + return ( + <> + + + + + + + {t("recommendations-to-others")} + + + + + + + + {/* Search Bar */} + {streamers.length < 8 && ( + + + + + + + + + + {searching && ( + <> + + + + Searching... + + + + )} + + {!searching && searchResults.length > 0 && ( + <> + + {searchResults.map((actor, index) => { + const alreadyAdded = streamers.includes(actor.did); + return ( + + {index > 0 && } + + !alreadyAdded && handleSelectActor(actor) + } + disabled={alreadyAdded} + > + {({ pressed }) => ( + + + @{actor.handle} + + {alreadyAdded && ( + + Added + + )} + + )} + + + ); + })} + + )} + + {!searching && + searchQuery.trim() && + searchResults.length === 0 && ( + <> + + + + No results found + + + + )} + + + {searchQuery.trim() === "" && ( + + )} + + )} + + {loading ? ( + + ) : ( + + + {streamers.length === 0 ? ( + + + {t("no-recommendations-yet")} + + + ) : ( + streamers.map((streamer, index) => ( + + {index > 0 && } + + + #{index + 1} + + + + handleStreamerChange(index, value) + } + onBlur={() => handleStreamerBlur(index)} + placeholder="did:plc:..." + /> + {errors[index] && ( + + {errors[index]} + + )} + + handleDelete(index)} + style={({ pressed }) => [ + { + padding: 8, + borderRadius: 6, + backgroundColor: pressed + ? "#ffffff08" + : "transparent", + }, + ]} + > + + + + + )) + )} + + {streamers.length > 0 && streamers.length < 8 && ( + + )} + + {streamers.length < 8 && ( + + {({ pressed }) => ( + + + Add DID manually + + )} + + )} + + + {saving && ( + + + {t("saving")} + + + )} + + )} + + + + + + !open && setDeleteDialog({ isVisible: false, index: null }) + } + title={t("delete")} + dismissible={false} + > + + {t("confirm-delete")} + + {t("action-cannot-be-undone")} + + + + + + + + + + ); +} diff --git a/js/app/components/settings/streaming-category-settings.tsx b/js/app/components/settings/streaming-category-settings.tsx index e9870e4d4..02bfdd1fb 100644 --- a/js/app/components/settings/streaming-category-settings.tsx +++ b/js/app/components/settings/streaming-category-settings.tsx @@ -5,7 +5,7 @@ import { View, zero, } from "@streamplace/components"; -import { Key, Webhook } from "lucide-react-native"; +import { Heart, Key, Webhook } from "lucide-react-native"; import { useTranslation } from "react-i18next"; import { ScrollView } from "react-native"; import { SettingsNavigationItem } from "./components/settings-navigation-item"; @@ -24,6 +24,12 @@ export function StreamingCategorySettings() { icon={Key} /> + + = { AccountCategory: "settings/account", StreamingCategory: "settings/streaming", WebhooksSettings: "settings/streaming/webhooks", + RecommendationsSettings: "settings/streaming/recommendations", PrivacyCategory: "settings/privacy", DanmuCategory: "settings/danmu", AdvancedCategory: "settings/advanced", @@ -751,6 +754,11 @@ const SettingsStack = () => { component={WebhookManager} options={{ headerTitle: "Webhooks", title: "Webhooks" }} /> + *[other] { $count } keys } +## Recommendations +recommendations = Recommendations +manage-recommendations = Manage Recommendations +recommendations-to-others = Recommendations to Others +recommendations-description = Share up to 8 streamers you recommend to your viewers +no-recommendations-yet = No recommendations configured yet +add-recommendation = Add Recommendation +streamer-did = Streamer DID +recommendations-count = { $count -> + [one] { $count } recommendation + *[other] { $count } recommendations +} + ## Webhook Management webhooks = Webhooks webhook-integrations = Webhook Integrations diff --git a/js/docs/src/content/docs/lex-reference/live/place-stream-live-getrecommendations.md b/js/docs/src/content/docs/lex-reference/live/place-stream-live-getrecommendations.md new file mode 100644 index 000000000..7218a6c40 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/live/place-stream-live-getrecommendations.md @@ -0,0 +1,84 @@ +--- +title: place.stream.live.getRecommendations +description: Reference for the place.stream.live.getRecommendations lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `query` + +Get the list of streamers recommended by a user + +**Parameters:** + +| Name | Type | Req'd | Description | Constraints | +| --------- | -------- | ----- | -------------------------------------------------- | ------------- | +| `userDID` | `string` | ✅ | The DID of the user whose recommendations to fetch | Format: `did` | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ----------- | ----------------- | ----- | --------------------------------------------- | ------------- | +| `streamers` | Array of `string` | ✅ | Ordered list of recommended streamer DIDs | | +| `userDID` | `string` | ❌ | The user who created this recommendation list | Format: `did` | + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.live.getRecommendations", + "defs": { + "main": { + "type": "query", + "description": "Get the list of streamers recommended by a user", + "parameters": { + "type": "params", + "required": ["userDID"], + "properties": { + "userDID": { + "type": "string", + "format": "did", + "description": "The DID of the user whose recommendations to fetch" + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["streamers"], + "properties": { + "streamers": { + "type": "array", + "description": "Ordered list of recommended streamer DIDs", + "items": { + "type": "string", + "format": "did" + } + }, + "userDID": { + "type": "string", + "format": "did", + "description": "The user who created this recommendation list" + } + } + } + } + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/live/place-stream-live-recommendations.md b/js/docs/src/content/docs/lex-reference/live/place-stream-live-recommendations.md new file mode 100644 index 000000000..bb70004cd --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/live/place-stream-live-recommendations.md @@ -0,0 +1,64 @@ +--- +title: place.stream.live.recommendations +description: Reference for the place.stream.live.recommendations lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `record` + +A list of recommended streamers, in order of preference + +**Record Key:** `self` + +**Record Properties:** + +| Name | Type | Req'd | Description | Constraints | +| ----------- | ----------------- | ----- | ----------------------------------------------------- | ----------------------------- | +| `streamers` | Array of `string` | ✅ | Ordered list of recommended streamer DIDs | Min Items: 0
Max Items: 8 | +| `createdAt` | `string` | ✅ | Client-declared timestamp when this list was created. | Format: `datetime` | + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.live.recommendations", + "defs": { + "main": { + "type": "record", + "description": "A list of recommended streamers, in order of preference", + "key": "self", + "record": { + "type": "object", + "required": ["streamers", "createdAt"], + "properties": { + "streamers": { + "type": "array", + "description": "Ordered list of recommended streamer DIDs", + "items": { + "type": "string", + "format": "did" + }, + "maxLength": 8, + "minLength": 0 + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "Client-declared timestamp when this list was created." + } + } + } + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/live/place-stream-live-searchactorstypeahead.md b/js/docs/src/content/docs/lex-reference/live/place-stream-live-searchactorstypeahead.md new file mode 100644 index 000000000..23ad983dc --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/live/place-stream-live-searchactorstypeahead.md @@ -0,0 +1,114 @@ +--- +title: place.stream.live.searchActorsTypeahead +description: Reference for the place.stream.live.searchActorsTypeahead lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `query` + +Find actor suggestions for a prefix search term. Expected use is for +auto-completion during text field entry. + +**Parameters:** + +| Name | Type | Req'd | Description | Constraints | +| ------- | --------- | ----- | --------------------------------------------- | ------------------------------------- | +| `q` | `string` | ❌ | Search query prefix; not a full query string. | | +| `limit` | `integer` | ❌ | | Min: 1
Max: 100
Default: `10` | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| -------- | --------------------------- | ----- | ----------- | ----------- | +| `actors` | Array of [`#actor`](#actor) | ✅ | | | + +--- + + + +### `actor` + +**Type:** `object` + +**Properties:** + +| Name | Type | Req'd | Description | Constraints | +| -------- | -------- | ----- | ------------------ | ---------------- | +| `did` | `string` | ✅ | The actor's DID | Format: `did` | +| `handle` | `string` | ✅ | The actor's handle | Format: `handle` | + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.live.searchActorsTypeahead", + "defs": { + "main": { + "type": "query", + "description": "Find actor suggestions for a prefix search term. Expected use is for auto-completion during text field entry.", + "parameters": { + "type": "params", + "properties": { + "q": { + "type": "string", + "description": "Search query prefix; not a full query string." + }, + "limit": { + "type": "integer", + "minimum": 1, + "maximum": 100, + "default": 10 + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["actors"], + "properties": { + "actors": { + "type": "array", + "items": { + "type": "ref", + "ref": "#actor" + } + } + } + } + } + }, + "actor": { + "type": "object", + "required": ["did", "handle"], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "The actor's DID" + }, + "handle": { + "type": "string", + "format": "handle", + "description": "The actor's handle" + } + } + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/openapi.json b/js/docs/src/content/docs/lex-reference/openapi.json index 4e13d6b69..f3dc74790 100644 --- a/js/docs/src/content/docs/lex-reference/openapi.json +++ b/js/docs/src/content/docs/lex-reference/openapi.json @@ -619,6 +619,54 @@ ] } }, + "/xrpc/place.stream.live.getRecommendations": { + "get": { + "summary": "Get the list of streamers recommended by a user", + "operationId": "place.stream.live.getRecommendations", + "tags": ["place.stream.live"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "streamers": { + "type": "array", + "description": "Ordered list of recommended streamer DIDs", + "items": { + "type": "string", + "format": "did" + } + }, + "userDID": { + "type": "string", + "description": "The user who created this recommendation list", + "format": "did" + } + }, + "required": ["streamers"] + } + } + } + } + }, + "parameters": [ + { + "name": "userDID", + "in": "query", + "required": true, + "description": "The DID of the user whose recommendations to fetch", + "schema": { + "type": "string", + "description": "The DID of the user whose recommendations to fetch", + "format": "did" + } + } + ] + } + }, "/xrpc/place.stream.live.getSegments": { "get": { "summary": "Get a list of livestream segments for a user", @@ -679,6 +727,57 @@ ] } }, + "/xrpc/place.stream.live.searchActorsTypeahead": { + "get": { + "summary": "Find actor suggestions for a prefix search term. Expected use is for auto-completion during text field entry.", + "operationId": "place.stream.live.searchActorsTypeahead", + "tags": ["place.stream.live"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "actors": { + "type": "array", + "items": { + "$ref": "#/components/schemas/place.stream.live.searchActorsTypeahead_actor" + } + } + }, + "required": ["actors"] + } + } + } + } + }, + "parameters": [ + { + "name": "q", + "in": "query", + "required": false, + "description": "Search query prefix; not a full query string.", + "schema": { + "type": "string", + "description": "Search query prefix; not a full query string." + } + }, + { + "name": "limit", + "in": "query", + "required": false, + "schema": { + "type": "integer", + "default": 10, + "minimum": 1, + "maximum": 100 + } + } + ] + } + }, "/xrpc/place.stream.live.subscribeSegments": { "get": { "summary": "Subscribe to a stream's new segments as they come in!", @@ -2155,6 +2254,22 @@ }, "required": ["cid", "record"] }, + "place.stream.live.searchActorsTypeahead_actor": { + "type": "object", + "properties": { + "did": { + "type": "string", + "description": "The actor's DID", + "format": "did" + }, + "handle": { + "type": "string", + "description": "The actor's handle", + "format": "handle" + } + }, + "required": ["did", "handle"] + }, "place.stream.live.subscribeSegments_segment": { "type": "string", "format": "byte", diff --git a/lexicons/place/stream/live/getRecommendations.json b/lexicons/place/stream/live/getRecommendations.json new file mode 100644 index 000000000..98fbed98d --- /dev/null +++ b/lexicons/place/stream/live/getRecommendations.json @@ -0,0 +1,43 @@ +{ + "lexicon": 1, + "id": "place.stream.live.getRecommendations", + "defs": { + "main": { + "type": "query", + "description": "Get the list of streamers recommended by a user", + "parameters": { + "type": "params", + "required": ["userDID"], + "properties": { + "userDID": { + "type": "string", + "format": "did", + "description": "The DID of the user whose recommendations to fetch" + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["streamers"], + "properties": { + "streamers": { + "type": "array", + "description": "Ordered list of recommended streamer DIDs", + "items": { + "type": "string", + "format": "did" + } + }, + "userDID": { + "type": "string", + "format": "did", + "description": "The user who created this recommendation list" + } + } + } + } + } + } +} diff --git a/lexicons/place/stream/live/recommendations.json b/lexicons/place/stream/live/recommendations.json new file mode 100644 index 000000000..babe8b068 --- /dev/null +++ b/lexicons/place/stream/live/recommendations.json @@ -0,0 +1,32 @@ +{ + "lexicon": 1, + "id": "place.stream.live.recommendations", + "defs": { + "main": { + "type": "record", + "description": "A list of recommended streamers, in order of preference", + "key": "self", + "record": { + "type": "object", + "required": ["streamers", "createdAt"], + "properties": { + "streamers": { + "type": "array", + "description": "Ordered list of recommended streamer DIDs", + "items": { + "type": "string", + "format": "did" + }, + "maxLength": 8, + "minLength": 0 + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "Client-declared timestamp when this list was created." + } + } + } + } + } +} diff --git a/lexicons/place/stream/live/searchActorsTypeahead.json b/lexicons/place/stream/live/searchActorsTypeahead.json new file mode 100644 index 000000000..46658f283 --- /dev/null +++ b/lexicons/place/stream/live/searchActorsTypeahead.json @@ -0,0 +1,57 @@ +{ + "lexicon": 1, + "id": "place.stream.live.searchActorsTypeahead", + "defs": { + "main": { + "type": "query", + "description": "Find actor suggestions for a prefix search term. Expected use is for auto-completion during text field entry.", + "parameters": { + "type": "params", + "properties": { + "q": { + "type": "string", + "description": "Search query prefix; not a full query string." + }, + "limit": { + "type": "integer", + "minimum": 1, + "maximum": 100, + "default": 10 + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["actors"], + "properties": { + "actors": { + "type": "array", + "items": { + "type": "ref", + "ref": "#actor" + } + } + } + } + } + }, + "actor": { + "type": "object", + "required": ["did", "handle"], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "The actor's DID" + }, + "handle": { + "type": "string", + "format": "handle", + "description": "The actor's handle" + } + } + } + } +} diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index 16c283c99..4133a4700 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -161,6 +161,7 @@ var CollectionFilter = []string{ constants.APP_BSKY_GRAPH_FOLLOW, constants.APP_BSKY_FEED_POST, constants.APP_BSKY_GRAPH_BLOCK, + "place.stream.live.recommendations", } func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt *comatproto.SyncSubscribeRepos_Commit) { diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index 3c5cd5d57..d324159fd 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -2,6 +2,7 @@ package atproto import ( "context" + "encoding/json" "errors" "fmt" "reflect" @@ -423,6 +424,38 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD log.Error(ctx, "failed to create metadata configuration", "err", err) } + case *streamplace.LiveRecommendations: + log.Debug(ctx, "creating recommendations", "userDID", userDID, "count", len(rec.Streamers)) + + // Validate max 8 streamers + if len(rec.Streamers) > 8 { + log.Warn(ctx, "recommendations exceed maximum of 8", "count", len(rec.Streamers)) + return fmt.Errorf("maximum 8 recommendations allowed, got %d", len(rec.Streamers)) + } + + // Marshal streamers to JSON + streamersJSON, err := json.Marshal(rec.Streamers) + if err != nil { + return fmt.Errorf("failed to marshal streamers: %w", err) + } + + // Parse createdAt timestamp + createdAt, err := time.Parse(time.RFC3339, rec.CreatedAt) + if err != nil { + return fmt.Errorf("failed to parse createdAt: %w", err) + } + + recommendation := &statedb.Recommendation{ + UserDID: userDID, + Streamers: json.RawMessage(streamersJSON), + CreatedAt: createdAt, + } + + err = atsync.StatefulDB.UpsertRecommendation(recommendation) + if err != nil { + return fmt.Errorf("failed to upsert recommendation: %w", err) + } + default: log.Debug(ctx, "unhandled record type", "type", reflect.TypeOf(rec)) } diff --git a/pkg/gen/gen.go b/pkg/gen/gen.go index 7ec89429f..821d580b2 100644 --- a/pkg/gen/gen.go +++ b/pkg/gen/gen.go @@ -32,6 +32,7 @@ func main() { streamplace.MetadataDistributionPolicy{}, streamplace.MetadataContentRights{}, streamplace.MetadataContentWarnings{}, + streamplace.LiveRecommendations{}, ); err != nil { panic(err) } diff --git a/pkg/model/model.go b/pkg/model/model.go index 8b100899d..db329fd87 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -48,6 +48,7 @@ type Model interface { GetRepoByHandleOrDID(arg string) (*Repo, error) GetRepoBySigningKey(signingKey string) (*Repo, error) GetAllRepos() ([]Repo, error) + SearchReposByHandle(query string, limit int) ([]Repo, error) UpdateRepo(repo *Repo) error UpdateSigningKey(key *SigningKey) error diff --git a/pkg/model/repo.go b/pkg/model/repo.go index 64ca807ce..44735a9c3 100644 --- a/pkg/model/repo.go +++ b/pkg/model/repo.go @@ -77,3 +77,14 @@ func (m *DBModel) GetRepoByHandleOrDID(arg string) (*Repo, error) { func (m *DBModel) UpdateRepo(repo *Repo) error { return m.DB.Save(repo).Error } + +func (m *DBModel) SearchReposByHandle(query string, limit int) ([]Repo, error) { + var repos []Repo + // Search for repos where handle starts with the query (case-insensitive) + // Use LIKE with LOWER for sqlite/postgres compatibility + res := m.DB.Where("LOWER(handle) LIKE LOWER(?)", query+"%").Limit(limit).Find(&repos) + if res.Error != nil { + return nil, res.Error + } + return repos, nil +} diff --git a/pkg/spxrpc/place_stream_live.go b/pkg/spxrpc/place_stream_live.go index 42f17e627..67c23a826 100644 --- a/pkg/spxrpc/place_stream_live.go +++ b/pkg/spxrpc/place_stream_live.go @@ -153,3 +153,28 @@ func (s *Server) handlePlaceStreamLiveSubscribeSegments(c echo.Context) error { log.Debug(c.Request().Context(), "received message", "message", string(msg)) } } + +func (s *Server) handlePlaceStreamLiveGetRecommendations(ctx context.Context, userDID string) (*placestreamtypes.LiveGetRecommendations_Output, error) { + if userDID == "" { + return nil, echo.NewHTTPError(http.StatusBadRequest, "userDID is required") + } + + rec, err := s.statefulDB.GetRecommendation(userDID) + if err != nil { + // If not found, return empty array + return &placestreamtypes.LiveGetRecommendations_Output{ + Streamers: []string{}, + UserDID: &userDID, + }, nil + } + + streamers, err := rec.GetStreamersArray() + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to parse recommendations") + } + + return &placestreamtypes.LiveGetRecommendations_Output{ + Streamers: streamers, + UserDID: &userDID, + }, nil +} diff --git a/pkg/spxrpc/place_stream_live_searchActorsTypeahead.go b/pkg/spxrpc/place_stream_live_searchActorsTypeahead.go new file mode 100644 index 000000000..5ff210e11 --- /dev/null +++ b/pkg/spxrpc/place_stream_live_searchActorsTypeahead.go @@ -0,0 +1,45 @@ +package spxrpc + +import ( + "context" + "net/http" + + "github.com/labstack/echo/v4" + placestreamtypes "stream.place/streamplace/pkg/streamplace" +) + +func (s *Server) handlePlaceStreamLiveSearchActorsTypeahead(ctx context.Context, limit int, q string) (*placestreamtypes.LiveSearchActorsTypeahead_Output, error) { + if q == "" { + return &placestreamtypes.LiveSearchActorsTypeahead_Output{ + Actors: []*placestreamtypes.LiveSearchActorsTypeahead_Actor{}, + }, nil + } + + // Default limit to 10 if not specified + searchLimit := 10 + if limit > 0 { + searchLimit = limit + if searchLimit > 100 { + searchLimit = 100 + } + } + + // Search repos by handle + repos, err := s.model.SearchReposByHandle(q, searchLimit) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to search actors "+err.Error()) + } + + // Convert to output format + actors := make([]*placestreamtypes.LiveSearchActorsTypeahead_Actor, len(repos)) + for i, repo := range repos { + actors[i] = &placestreamtypes.LiveSearchActorsTypeahead_Actor{ + Did: repo.DID, + Handle: repo.Handle, + } + } + + return &placestreamtypes.LiveSearchActorsTypeahead_Output{ + Actors: actors, + }, nil +} diff --git a/pkg/spxrpc/stubs.go b/pkg/spxrpc/stubs.go index 81718267c..70ab7b475 100644 --- a/pkg/spxrpc/stubs.go +++ b/pkg/spxrpc/stubs.go @@ -266,7 +266,9 @@ func (s *Server) RegisterHandlersPlaceStream(e *echo.Echo) error { e.GET("/xrpc/place.stream.graph.getFollowingUser", s.HandlePlaceStreamGraphGetFollowingUser) e.GET("/xrpc/place.stream.live.getLiveUsers", s.HandlePlaceStreamLiveGetLiveUsers) e.GET("/xrpc/place.stream.live.getProfileCard", s.HandlePlaceStreamLiveGetProfileCard) + e.GET("/xrpc/place.stream.live.getRecommendations", s.HandlePlaceStreamLiveGetRecommendations) e.GET("/xrpc/place.stream.live.getSegments", s.HandlePlaceStreamLiveGetSegments) + e.GET("/xrpc/place.stream.live.searchActorsTypeahead", s.HandlePlaceStreamLiveSearchActorsTypeahead) e.POST("/xrpc/place.stream.server.createWebhook", s.HandlePlaceStreamServerCreateWebhook) e.POST("/xrpc/place.stream.server.deleteWebhook", s.HandlePlaceStreamServerDeleteWebhook) e.GET("/xrpc/place.stream.server.getServerTime", s.HandlePlaceStreamServerGetServerTime) @@ -343,6 +345,20 @@ func (s *Server) HandlePlaceStreamLiveGetProfileCard(c echo.Context) error { return c.Stream(200, "application/octet-stream", out) } +func (s *Server) HandlePlaceStreamLiveGetRecommendations(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamLiveGetRecommendations") + defer span.End() + userDID := c.QueryParam("userDID") + var out *placestreamtypes.LiveGetRecommendations_Output + var handleErr error + // func (s *Server) handlePlaceStreamLiveGetRecommendations(ctx context.Context,userDID string) (*placestreamtypes.LiveGetRecommendations_Output, error) + out, handleErr = s.handlePlaceStreamLiveGetRecommendations(ctx, userDID) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + func (s *Server) HandlePlaceStreamLiveGetSegments(c echo.Context) error { ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamLiveGetSegments") defer span.End() @@ -369,6 +385,31 @@ func (s *Server) HandlePlaceStreamLiveGetSegments(c echo.Context) error { return c.JSON(200, out) } +func (s *Server) HandlePlaceStreamLiveSearchActorsTypeahead(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamLiveSearchActorsTypeahead") + defer span.End() + + var limit int + if p := c.QueryParam("limit"); p != "" { + var err error + limit, err = strconv.Atoi(p) + if err != nil { + return err + } + } else { + limit = 10 + } + q := c.QueryParam("q") + var out *placestreamtypes.LiveSearchActorsTypeahead_Output + var handleErr error + // func (s *Server) handlePlaceStreamLiveSearchActorsTypeahead(ctx context.Context,limit int,q string) (*placestreamtypes.LiveSearchActorsTypeahead_Output, error) + out, handleErr = s.handlePlaceStreamLiveSearchActorsTypeahead(ctx, limit, q) + if handleErr != nil { + return handleErr + } + 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() diff --git a/pkg/statedb/recommendations.go b/pkg/statedb/recommendations.go new file mode 100644 index 000000000..1f35ead86 --- /dev/null +++ b/pkg/statedb/recommendations.go @@ -0,0 +1,93 @@ +package statedb + +import ( + "encoding/json" + "errors" + "fmt" + "time" + + "gorm.io/gorm" +) + +type Recommendation struct { + UserDID string `gorm:"column:user_did;primaryKey"` + Streamers json.RawMessage `gorm:"column:streamers;type:json;not null"` + CreatedAt time.Time `gorm:"column:created_at"` + UpdatedAt time.Time `gorm:"column:updated_at"` +} + +func (r *Recommendation) TableName() string { + return "recommendations" +} + +// UpsertRecommendation creates or updates recommendations for a user +func (state *StatefulDB) UpsertRecommendation(rec *Recommendation) error { + if rec.UserDID == "" { + return fmt.Errorf("user DID cannot be empty") + } + + // Validate JSON contains array of max 8 DIDs + var streamers []string + if err := json.Unmarshal(rec.Streamers, &streamers); err != nil { + return fmt.Errorf("invalid streamers JSON: %w", err) + } + if len(streamers) > 8 { + return fmt.Errorf("maximum 8 recommendations allowed, got %d", len(streamers)) + } + + now := time.Now() + if rec.CreatedAt.IsZero() { + rec.CreatedAt = now + } + rec.UpdatedAt = now + + // Use GORM's upsert (On Conflict Do Update) + result := state.DB.Save(rec) + if result.Error != nil { + return fmt.Errorf("database upsert failed: %w", result.Error) + } + + return nil +} + +// GetRecommendation retrieves recommendations for a user +func (state *StatefulDB) GetRecommendation(userDID string) (*Recommendation, error) { + if userDID == "" { + return nil, fmt.Errorf("user DID cannot be empty") + } + + var rec Recommendation + err := state.DB.Where("user_did = ?", userDID).First(&rec).Error + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, err + } + return nil, fmt.Errorf("database query failed: %w", err) + } + return &rec, nil +} + +// DeleteRecommendation removes recommendations for a user +func (state *StatefulDB) DeleteRecommendation(userDID string) error { + if userDID == "" { + return fmt.Errorf("user DID cannot be empty") + } + + result := state.DB.Where("user_did = ?", userDID).Delete(&Recommendation{}) + if result.Error != nil { + return fmt.Errorf("database delete failed: %w", result.Error) + } + if result.RowsAffected == 0 { + return fmt.Errorf("recommendation not found") + } + return nil +} + +// GetStreamersArray is a helper to unmarshal the streamers JSON into a slice +func (r *Recommendation) GetStreamersArray() ([]string, error) { + var streamers []string + if err := json.Unmarshal(r.Streamers, &streamers); err != nil { + return nil, fmt.Errorf("failed to unmarshal streamers: %w", err) + } + return streamers, nil +} diff --git a/pkg/statedb/statedb.go b/pkg/statedb/statedb.go index e319deac1..b1410cb21 100644 --- a/pkg/statedb/statedb.go +++ b/pkg/statedb/statedb.go @@ -49,6 +49,7 @@ var StatefulDBModels = []any{ AppTask{}, Repo{}, Webhook{}, + Recommendation{}, } var NoPostgresDatabaseCode = "3D000" diff --git a/pkg/streamplace/cbor_gen.go b/pkg/streamplace/cbor_gen.go index b58bdda87..d434d7573 100644 --- a/pkg/streamplace/cbor_gen.go +++ b/pkg/streamplace/cbor_gen.go @@ -4975,3 +4975,206 @@ func (t *MetadataContentWarnings) UnmarshalCBOR(r io.Reader) (err error) { return nil } +func (t *LiveRecommendations) MarshalCBOR(w io.Writer) error { + if t == nil { + _, err := w.Write(cbg.CborNull) + return err + } + + cw := cbg.NewCborWriter(w) + + if _, err := cw.Write([]byte{163}); err != nil { + return err + } + + // t.LexiconTypeID (string) (string) + if len("$type") > 1000000 { + return xerrors.Errorf("Value in field \"$type\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("$type"))); err != nil { + return err + } + if _, err := cw.WriteString(string("$type")); err != nil { + return err + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("place.stream.live.recommendations"))); err != nil { + return err + } + if _, err := cw.WriteString(string("place.stream.live.recommendations")); err != nil { + return err + } + + // t.CreatedAt (string) (string) + if len("createdAt") > 1000000 { + return xerrors.Errorf("Value in field \"createdAt\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("createdAt"))); err != nil { + return err + } + if _, err := cw.WriteString(string("createdAt")); err != nil { + return err + } + + if len(t.CreatedAt) > 1000000 { + return xerrors.Errorf("Value in field t.CreatedAt was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.CreatedAt))); err != nil { + return err + } + if _, err := cw.WriteString(string(t.CreatedAt)); err != nil { + return err + } + + // t.Streamers ([]string) (slice) + if len("streamers") > 1000000 { + return xerrors.Errorf("Value in field \"streamers\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("streamers"))); err != nil { + return err + } + if _, err := cw.WriteString(string("streamers")); err != nil { + return err + } + + if len(t.Streamers) > 8192 { + return xerrors.Errorf("Slice value in field t.Streamers was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajArray, uint64(len(t.Streamers))); err != nil { + return err + } + for _, v := range t.Streamers { + if len(v) > 1000000 { + return xerrors.Errorf("Value in field v was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(v))); err != nil { + return err + } + if _, err := cw.WriteString(string(v)); err != nil { + return err + } + + } + return nil +} + +func (t *LiveRecommendations) UnmarshalCBOR(r io.Reader) (err error) { + *t = LiveRecommendations{} + + cr := cbg.NewCborReader(r) + + maj, extra, err := cr.ReadHeader() + if err != nil { + return err + } + defer func() { + if err == io.EOF { + err = io.ErrUnexpectedEOF + } + }() + + if maj != cbg.MajMap { + return fmt.Errorf("cbor input should be of type map") + } + + if extra > cbg.MaxLength { + return fmt.Errorf("LiveRecommendations: map struct too large (%d)", extra) + } + + n := extra + + nameBuf := make([]byte, 9) + for i := uint64(0); i < n; i++ { + nameLen, ok, err := cbg.ReadFullStringIntoBuf(cr, nameBuf, 1000000) + if err != nil { + return err + } + + if !ok { + // Field doesn't exist on this type, so ignore it + if err := cbg.ScanForLinks(cr, func(cid.Cid) {}); err != nil { + return err + } + continue + } + + switch string(nameBuf[:nameLen]) { + // t.LexiconTypeID (string) (string) + case "$type": + + { + sval, err := cbg.ReadStringWithMax(cr, 1000000) + if err != nil { + return err + } + + t.LexiconTypeID = string(sval) + } + // t.CreatedAt (string) (string) + case "createdAt": + + { + sval, err := cbg.ReadStringWithMax(cr, 1000000) + if err != nil { + return err + } + + t.CreatedAt = string(sval) + } + // t.Streamers ([]string) (slice) + case "streamers": + + maj, extra, err = cr.ReadHeader() + if err != nil { + return err + } + + if extra > 8192 { + return fmt.Errorf("t.Streamers: array too large (%d)", extra) + } + + if maj != cbg.MajArray { + return fmt.Errorf("expected cbor array") + } + + if extra > 0 { + t.Streamers = make([]string, extra) + } + + for i := 0; i < int(extra); i++ { + { + var maj byte + var extra uint64 + var err error + _ = maj + _ = extra + _ = err + + { + sval, err := cbg.ReadStringWithMax(cr, 1000000) + if err != nil { + return err + } + + t.Streamers[i] = string(sval) + } + + } + } + + default: + // Field doesn't exist on this type, so ignore it + if err := cbg.ScanForLinks(r, func(cid.Cid) {}); err != nil { + return err + } + } + } + + return nil +} diff --git a/pkg/streamplace/livegetRecommendations.go b/pkg/streamplace/livegetRecommendations.go new file mode 100644 index 000000000..dab1dbef0 --- /dev/null +++ b/pkg/streamplace/livegetRecommendations.go @@ -0,0 +1,34 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package streamplace + +// schema: place.stream.live.getRecommendations + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +// LiveGetRecommendations_Output is the output of a place.stream.live.getRecommendations call. +type LiveGetRecommendations_Output struct { + // streamers: Ordered list of recommended streamer DIDs + Streamers []string `json:"streamers" cborgen:"streamers"` + // userDID: The user who created this recommendation list + UserDID *string `json:"userDID,omitempty" cborgen:"userDID,omitempty"` +} + +// LiveGetRecommendations calls the XRPC method "place.stream.live.getRecommendations". +// +// userDID: The DID of the user whose recommendations to fetch +func LiveGetRecommendations(ctx context.Context, c util.LexClient, userDID string) (*LiveGetRecommendations_Output, error) { + var out LiveGetRecommendations_Output + + params := map[string]interface{}{} + params["userDID"] = userDID + if err := c.LexDo(ctx, util.Query, "", "place.stream.live.getRecommendations", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/pkg/streamplace/liverecommendations.go b/pkg/streamplace/liverecommendations.go new file mode 100644 index 000000000..7a59cdb84 --- /dev/null +++ b/pkg/streamplace/liverecommendations.go @@ -0,0 +1,21 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package streamplace + +// schema: place.stream.live.recommendations + +import ( + "github.com/bluesky-social/indigo/lex/util" +) + +func init() { + util.RegisterType("place.stream.live.recommendations", &LiveRecommendations{}) +} // +// RECORDTYPE: LiveRecommendations +type LiveRecommendations struct { + LexiconTypeID string `json:"$type,const=place.stream.live.recommendations" cborgen:"$type,const=place.stream.live.recommendations"` + // createdAt: Client-declared timestamp when this list was created. + CreatedAt string `json:"createdAt" cborgen:"createdAt"` + // streamers: Ordered list of recommended streamer DIDs + Streamers []string `json:"streamers" cborgen:"streamers"` +} diff --git a/pkg/streamplace/livesearchActorsTypeahead.go b/pkg/streamplace/livesearchActorsTypeahead.go new file mode 100644 index 000000000..864fbedc4 --- /dev/null +++ b/pkg/streamplace/livesearchActorsTypeahead.go @@ -0,0 +1,44 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package streamplace + +// schema: place.stream.live.searchActorsTypeahead + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +// LiveSearchActorsTypeahead_Actor is a "actor" in the place.stream.live.searchActorsTypeahead schema. +type LiveSearchActorsTypeahead_Actor struct { + // did: The actor's DID + Did string `json:"did" cborgen:"did"` + // handle: The actor's handle + Handle string `json:"handle" cborgen:"handle"` +} + +// LiveSearchActorsTypeahead_Output is the output of a place.stream.live.searchActorsTypeahead call. +type LiveSearchActorsTypeahead_Output struct { + Actors []*LiveSearchActorsTypeahead_Actor `json:"actors" cborgen:"actors"` +} + +// LiveSearchActorsTypeahead calls the XRPC method "place.stream.live.searchActorsTypeahead". +// +// q: Search query prefix; not a full query string. +func LiveSearchActorsTypeahead(ctx context.Context, c util.LexClient, limit int64, q string) (*LiveSearchActorsTypeahead_Output, error) { + var out LiveSearchActorsTypeahead_Output + + params := map[string]interface{}{} + if limit != 0 { + params["limit"] = limit + } + if q != "" { + params["q"] = q + } + if err := c.LexDo(ctx, util.Query, "", "place.stream.live.searchActorsTypeahead", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +}