diff --git a/Makefile b/Makefile index 0c2d62e4..61e1a0ba 100644 --- a/Makefile +++ b/Makefile @@ -384,6 +384,7 @@ js-lexicons: && sed -i.bak 's/AppBskyGraphBlock\.Main/AppBskyGraphBlock\.Record/' $$(find ./js/streamplace/src/lexicons/types/place/stream -type f) \ && sed -i.bak 's/PlaceStreamMultistreamTarget\.Main/PlaceStreamMultistreamTarget\.Record/' $$(find ./js/streamplace/src/lexicons/types/place/stream -type f) \ && sed -i.bak 's/PlaceStreamChatProfile\.Main/PlaceStreamChatProfile\.Record/' $$(find ./js/streamplace/src/lexicons/types/place/stream -type f) \ + && sed -i.bak 's/PlaceStreamLivestream\.Main/PlaceStreamLivestream\.Record/' $$(find ./js/streamplace/src/lexicons/types/place/stream/live -type f) \ && for x in $$(find ./js/streamplace/src/lexicons -type f -name '*.ts'); do \ echo 'import { ComAtprotoSyncGetRepo, AppBskyRichtextFacet, AppBskyGraphBlock, ComAtprotoRepoStrongRef, AppBskyActorDefs, ComAtprotoSyncListRepos, AppBskyActorGetProfile, AppBskyFeedGetFeedSkeleton, ComAtprotoIdentityResolveHandle, ComAtprotoModerationCreateReport, ComAtprotoRepoCreateRecord, ComAtprotoRepoDeleteRecord, ComAtprotoRepoDescribeRepo, ComAtprotoRepoGetRecord, ComAtprotoRepoListRecords, ComAtprotoRepoPutRecord, ComAtprotoRepoUploadBlob, ComAtprotoServerDescribeServer, ComAtprotoSyncGetRecord, ComAtprotoSyncListReposComAtprotoRepoCreateRecord, ComAtprotoRepoDeleteRecord, ComAtprotoRepoGetRecord, ComAtprotoRepoListRecords, ComAtprotoIdentityRefreshIdentity } from "@atproto/api"' >> $$x; \ done \ diff --git a/js/app/components/live-dashboard/livestream-panel.tsx b/js/app/components/live-dashboard/livestream-panel.tsx index f733fc24..cf70349a 100644 --- a/js/app/components/live-dashboard/livestream-panel.tsx +++ b/js/app/components/live-dashboard/livestream-panel.tsx @@ -353,7 +353,7 @@ function LivestreamPanel({ scrollable = true }: { scrollable?: boolean }) { if (!livestream) return; setEndingLivestream(true); try { - await endLivestream(livestream); + await endLivestream(); } catch (error) { console.error("Error ending livestream:", error); toast.show("Error", "Failed to end livestream", { diff --git a/js/components/src/components/mobile-player/video.tsx b/js/components/src/components/mobile-player/video.tsx index fb57b1ed..085b1ddf 100644 --- a/js/components/src/components/mobile-player/video.tsx +++ b/js/components/src/components/mobile-player/video.tsx @@ -544,6 +544,7 @@ export function WebRTCPlayerInner({ export function WebcamIngestPlayer(props: VideoProps) { const ingestMediaSource = usePlayerStore((x) => x.ingestMediaSource); const ingestAutoStart = usePlayerStore((x) => x.ingestAutoStart); + const setIngestLive = usePlayerStore((x) => x.setIngestLive); const [error, setError] = useState(null); @@ -606,15 +607,17 @@ export function WebcamIngestPlayer(props: VideoProps) { }, [ingestMediaSource]); useEffect(() => { - if (!ingestAutoStart) { - setRemoteMediaStream(null); - return; - } + // if (!ingestAutoStart) { + // setRemoteMediaStream(null); + // return; + // } if (!localMediaStream) { return; } + console.log("setting remote media stream", localMediaStream); + setIngestLive(true); setRemoteMediaStream(localMediaStream); - }, [localMediaStream, ingestAutoStart]); + }, [localMediaStream, setIngestLive, setRemoteMediaStream]); useEffect(() => { if (!videoElement) { diff --git a/js/components/src/hooks/useLivestreamInfo.ts b/js/components/src/hooks/useLivestreamInfo.ts index 3bef6835..77dd9688 100644 --- a/js/components/src/hooks/useLivestreamInfo.ts +++ b/js/components/src/hooks/useLivestreamInfo.ts @@ -1,13 +1,12 @@ import { useState } from "react"; import { useLivestreamStore } from "../livestream-store"; import { usePlayerStore } from "../player-store"; -import { useCreateStreamRecord } from "../streamplace-store"; +import { useCreateStreamRecord, useEndLivestream } from "../streamplace-store"; export function useLivestreamInfo(url?: string) { const ingest = usePlayerStore((x) => x.ingestConnectionState); const profile = useLivestreamStore((x) => x.profile); - const setIngestLive = usePlayerStore((x) => x.setIngestLive); - const stopIngest = usePlayerStore((x) => x.stopIngest); + const endLivestream = useEndLivestream(); const createStreamRecord = useCreateStreamRecord(); @@ -49,7 +48,7 @@ export function useLivestreamInfo(url?: string) { // Stop the current broadcast const toggleStopStream = () => { console.log("Stopping stream..."); - stopIngest(); + endLivestream(); }; return { diff --git a/js/components/src/streamplace-store/stream.tsx b/js/components/src/streamplace-store/stream.tsx index a5c7d061..dceb13f1 100644 --- a/js/components/src/streamplace-store/stream.tsx +++ b/js/components/src/streamplace-store/stream.tsx @@ -127,7 +127,6 @@ export function useCreateStreamRecord() { let agent = usePDSAgent(); let url = useUrl(); const uploadThumbnail = useUploadThumbnail(); - return async ({ title, customThumbnail, @@ -143,97 +142,10 @@ export function useCreateStreamRecord() { notificationSettings?: PlaceStreamLivestream.NotificationSettings; idleTimeoutSeconds?: number; }) => { - if (typeof submitPost !== "boolean") { - submitPost = true; - } if (!agent) { throw new Error("No PDS agent found"); } - if (!agent.did) { - throw new Error("No user DID found, assuming not logged in"); - } - - const u = new URL(url); - - // let thumbnail: BlobRef | undefined = undefined; - - // if (customThumbnail) { - // try { - // thumbnail = await uploadThumbnail(agent, customThumbnail); - // } catch (e) { - // throw new Error(`Custom thumbnail upload failed ${e}`); - // } - // } else { - // // No custom thumbnail: fetch the server-side image and upload it - // // try thrice lel - // let tries = 0; - // try { - // for (; tries < 3; tries++) { - // try { - // console.log( - // `Fetching thumbnail from ${u.protocol}//${u.host}/api/playback/${agent.did}/stream.png`, - // ); - // const thumbnailRes = await fetch( - // `${u.protocol}//${u.host}/api/playback/${agent.did}/stream.png`, - // ); - // if (!thumbnailRes.ok) { - // throw new Error( - // `Failed to fetch thumbnail: ${thumbnailRes.status})`, - // ); - // } - // const thumbnailBlob = await thumbnailRes.blob(); - // console.log(thumbnailBlob); - // thumbnail = await uploadThumbnail(agent, thumbnailBlob); - // } catch (e) { - // console.warn( - // `Failed to fetch thumbnail, retrying (${tries + 1}/3): ${e}`, - // ); - // // Wait 1 second before retrying - // await new Promise((resolve) => setTimeout(resolve, 2000)); - // if (tries === 2) { - // throw new Error(`Failed to fetch thumbnail after 3 tries: ${e}`); - // } - // } - // } - // } catch (e) { - // throw new Error(`Thumbnail upload failed ${e}`); - // } - // } - - let newPost: undefined | { uri: string; cid: string } = undefined; - - const did = agent.did; - const profile = await agent.getProfile({ actor: did }); - - if (submitPost) { - if (!profile) { - throw new Error("No profile found for the user DID"); - } - - const params = new URLSearchParams({ - did: did, - time: new Date().toISOString(), - }); - - let post = await buildGoLivePost( - title, - u, - profile.data, - params, - undefined, - agent, - ); - - newPost = await createNewPost(agent, post); - - if (!newPost.uri || !newPost.cid) { - throw new Error( - "Cannot read properties of undefined (reading 'uri' or 'cid')", - ); - } - } - let platform: string = Platform.OS; let platVersion: string = Platform.Version ? Platform.Version.toString() @@ -246,12 +158,12 @@ export function useCreateStreamRecord() { ) { platVersion = getBrowserName(window.navigator.userAgent); } - - const thisUrl = `${url}/${profile.data.handle}`; - if (!canonicalUrl) { - canonicalUrl = thisUrl; + if (!agent.did) { + throw new Error("No user DID found, assuming not logged in"); } + const thisUrl = `${url}/${agent.did}`; + const record: PlaceStreamLivestream.Record = { $type: "place.stream.livestream", title: title, @@ -263,22 +175,29 @@ export function useCreateStreamRecord() { // user agent style string // e.g. `@streamplace/components/0.1.0 (ios, 32.0)` agent: `@streamplace/components/${PackageJson.version} (${platform}, ${platVersion})`, - post: newPost, - // thumb: thumbnail, idleTimeoutSeconds: idleTimeoutSeconds, }; - console.log("record", record); if (notificationSettings) { record.notificationSettings = notificationSettings; } - await agent.com.atproto.repo.createRecord({ - repo: agent.did, - collection: "place.stream.livestream", - record, + if (customThumbnail) { + try { + const thumbnail = await uploadThumbnail(agent, customThumbnail); + record.thumb = thumbnail; + } catch (e) { + throw new Error(`Custom thumbnail upload failed ${e}`); + } + } + + const output = await agent.place.stream.live.startLivestream({ + livestream: record, + streamer: agent.did, + createBlueskyPost: submitPost, }); - return record; + + return output; }; } @@ -347,7 +266,7 @@ export function useUpdateStreamRecord(customUrl: string | null = null) { export function useEndLivestream() { let agent = usePDSAgent(); - return async (livestream: LivestreamViewHydrated | null) => { + return async () => { if (!agent) { throw new Error("No PDS agent found"); } @@ -356,30 +275,6 @@ export function useEndLivestream() { throw new Error("No user DID found, assuming not logged in"); } - if (!livestream) { - throw new Error("No latest record"); - } - - let rkey = livestream.uri.split("/").pop(); - if (!rkey) { - throw new Error("No rkey?"); - } - - if (livestream.record.endedAt) { - throw new Error("Livestream already ended"); - } - - let record: PlaceStreamLivestream.Record = { - ...livestream.record, - endedAt: new Date().toISOString(), - }; - - await agent.com.atproto.repo.putRecord({ - repo: agent.did, - collection: "place.stream.livestream", - rkey, - record, - }); - return record; + return await agent.place.stream.live.stopLivestream({}); }; } diff --git a/js/docs/src/content/docs/lex-reference/live/place-stream-live-startlivestream.md b/js/docs/src/content/docs/lex-reference/live/place-stream-live-startlivestream.md new file mode 100644 index 00000000..e8a6df56 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/live/place-stream-live-startlivestream.md @@ -0,0 +1,101 @@ +--- +title: place.stream.live.startLivestream +description: Reference for the place.stream.live.startLivestream lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `procedure` + +Create a new place.stream.livestream record, automatically populating a thumbnail and creating a Bluesky post and whatnot. You can do this manually by creating a record but this method can work better for mobile livestreaming and such. + +**Parameters:** _(None defined)_ + +**Input:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ------------------- | ------------------------------------------------------------------- | ----- | ----------------------------------------------------------- | ------------- | +| `livestream` | [`place.stream.livestream`](/lex-reference/place-stream-livestream) | ✅ | | | +| `streamer` | `string` | ✅ | The DID of the streamer. | Format: `did` | +| `createBlueskyPost` | `boolean` | ❌ | Whether to create a Bluesky post announcing the livestream. | | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ----- | -------- | ----- | --------------------------------- | ------------- | +| `uri` | `string` | ✅ | The URI of the livestream record. | Format: `uri` | +| `cid` | `string` | ✅ | The CID of the livestream record. | Format: `cid` | + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.live.startLivestream", + "defs": { + "main": { + "type": "procedure", + "description": "Create a new place.stream.livestream record, automatically populating a thumbnail and creating a Bluesky post and whatnot. You can do this manually by creating a record but this method can work better for mobile livestreaming and such.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["streamer", "livestream"], + "properties": { + "livestream": { + "type": "ref", + "ref": "place.stream.livestream" + }, + "streamer": { + "type": "string", + "format": "did", + "description": "The DID of the streamer." + }, + "createBlueskyPost": { + "type": "boolean", + "description": "Whether to create a Bluesky post announcing the livestream." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["uri", "cid"], + "properties": { + "uri": { + "type": "string", + "format": "uri", + "description": "The URI of the livestream record." + }, + "cid": { + "type": "string", + "format": "cid", + "description": "The CID of the livestream record." + } + } + } + } + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/live/place-stream-live-stoplivestream.md b/js/docs/src/content/docs/lex-reference/live/place-stream-live-stoplivestream.md new file mode 100644 index 00000000..071a992e --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/live/place-stream-live-stoplivestream.md @@ -0,0 +1,82 @@ +--- +title: place.stream.live.stopLivestream +description: Reference for the place.stream.live.stopLivestream lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `procedure` + +Stop your current livestream, updating your current place.stream.livestream record and ceasing the flow of video. + +**Parameters:** _(None defined)_ + +**Input:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +_(No properties defined)_ +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ----- | -------- | ----- | --------------------------------------------- | ------------- | +| `uri` | `string` | ✅ | The URI of the stopped livestream record. | Format: `uri` | +| `cid` | `string` | ✅ | The new CID of the stopped livestream record. | Format: `cid` | + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.live.stopLivestream", + "defs": { + "main": { + "type": "procedure", + "description": "Stop your current livestream, updating your current place.stream.livestream record and ceasing the flow of video.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [], + "properties": {} + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["uri", "cid"], + "properties": { + "uri": { + "type": "string", + "format": "uri", + "description": "The URI of the stopped livestream record." + }, + "cid": { + "type": "string", + "format": "cid", + "description": "The new CID of the stopped livestream record." + } + } + } + } + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/openapi.json b/js/docs/src/content/docs/lex-reference/openapi.json index 44570b9f..41d6ddcb 100644 --- a/js/docs/src/content/docs/lex-reference/openapi.json +++ b/js/docs/src/content/docs/lex-reference/openapi.json @@ -1553,6 +1553,106 @@ ] } }, + "/xrpc/place.stream.live.startLivestream": { + "post": { + "summary": "Create a new place.stream.livestream record, automatically populating a thumbnail and creating a Bluesky post and whatnot. You can do this manually by creating a record but this method can work better for mobile livestreaming and such.", + "operationId": "place.stream.live.startLivestream", + "tags": ["place.stream.live"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "uri": { + "type": "string", + "description": "The URI of the livestream record.", + "format": "uri" + }, + "cid": { + "type": "string", + "description": "The CID of the livestream record.", + "format": "cid" + } + }, + "required": ["uri", "cid"] + } + } + } + } + }, + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "livestream": { + "$ref": "#/components/schemas/place.stream.livestream" + }, + "streamer": { + "type": "string", + "description": "The DID of the streamer.", + "format": "did" + }, + "createBlueskyPost": { + "type": "boolean", + "description": "Whether to create a Bluesky post announcing the livestream." + } + }, + "required": ["streamer", "livestream"] + } + } + } + } + } + }, + "/xrpc/place.stream.live.stopLivestream": { + "post": { + "summary": "Stop your current livestream, updating your current place.stream.livestream record and ceasing the flow of video.", + "operationId": "place.stream.live.stopLivestream", + "tags": ["place.stream.live"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "uri": { + "type": "string", + "description": "The URI of the stopped livestream record.", + "format": "uri" + }, + "cid": { + "type": "string", + "description": "The new CID of the stopped livestream record.", + "format": "cid" + } + }, + "required": ["uri", "cid"] + } + } + } + } + }, + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": {} + } + } + } + } + } + }, "/xrpc/place.stream.live.subscribeSegments": { "get": { "summary": "Subscribe to a stream's new segments as they come in!", @@ -3451,6 +3551,9 @@ }, "required": ["did", "handle"] }, + "place.stream.livestream": { + "description": "Unknown type" + }, "place.stream.live.subscribeSegments_segment": { "type": "string", "format": "byte", diff --git a/lexicons/place/stream/live/startLivestream.json b/lexicons/place/stream/live/startLivestream.json new file mode 100644 index 00000000..4509a0e9 --- /dev/null +++ b/lexicons/place/stream/live/startLivestream.json @@ -0,0 +1,51 @@ +{ + "lexicon": 1, + "id": "place.stream.live.startLivestream", + "defs": { + "main": { + "type": "procedure", + "description": "Create a new place.stream.livestream record, automatically populating a thumbnail and creating a Bluesky post and whatnot. You can do this manually by creating a record but this method can work better for mobile livestreaming and such.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["streamer", "livestream"], + "properties": { + "livestream": { + "type": "ref", + "ref": "place.stream.livestream" + }, + "streamer": { + "type": "string", + "format": "did", + "description": "The DID of the streamer." + }, + "createBlueskyPost": { + "type": "boolean", + "description": "Whether to create a Bluesky post announcing the livestream." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["uri", "cid"], + "properties": { + "uri": { + "type": "string", + "format": "uri", + "description": "The URI of the livestream record." + }, + "cid": { + "type": "string", + "format": "cid", + "description": "The CID of the livestream record." + } + } + } + } + } + } +} diff --git a/lexicons/place/stream/live/stopLivestream.json b/lexicons/place/stream/live/stopLivestream.json new file mode 100644 index 00000000..78ba2e16 --- /dev/null +++ b/lexicons/place/stream/live/stopLivestream.json @@ -0,0 +1,37 @@ +{ + "lexicon": 1, + "id": "place.stream.live.stopLivestream", + "defs": { + "main": { + "type": "procedure", + "description": "Stop your current livestream, updating your current place.stream.livestream record and ceasing the flow of video.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [], + "properties": {} + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["uri", "cid"], + "properties": { + "uri": { + "type": "string", + "format": "uri", + "description": "The URI of the stopped livestream record." + }, + "cid": { + "type": "string", + "format": "cid", + "description": "The new CID of the stopped livestream record." + } + } + } + } + } + } +} diff --git a/pkg/spxrpc/place_stream_live.go b/pkg/spxrpc/place_stream_live.go index 9e296bbd..6c19ac20 100644 --- a/pkg/spxrpc/place_stream_live.go +++ b/pkg/spxrpc/place_stream_live.go @@ -1,26 +1,34 @@ package spxrpc import ( + "bytes" "context" "encoding/json" "fmt" "net/http" "net/url" + "os" "strconv" "time" - "github.com/bluesky-social/indigo/lex/util" + comatproto "github.com/bluesky-social/indigo/api/atproto" + bsky "github.com/bluesky-social/indigo/api/bsky" + "github.com/bluesky-social/indigo/atproto/syntax" + lexutil "github.com/bluesky-social/indigo/lex/util" + "github.com/bluesky-social/indigo/util" + "github.com/bluesky-social/indigo/xrpc" "github.com/gorilla/websocket" "github.com/labstack/echo/v4" "github.com/streamplace/oatproxy/pkg/oatproxy" + "stream.place/streamplace/pkg/aqtime" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/spid" "stream.place/streamplace/pkg/spmetrics" - placestreamtypes "stream.place/streamplace/pkg/streamplace" + placestream "stream.place/streamplace/pkg/streamplace" ) -func (s *Server) handlePlaceStreamLiveDenyTeleport(ctx context.Context, input *placestreamtypes.LiveDenyTeleport_Input) (*placestreamtypes.LiveDenyTeleport_Output, error) { +func (s *Server) handlePlaceStreamLiveDenyTeleport(ctx context.Context, input *placestream.LiveDenyTeleport_Input) (*placestream.LiveDenyTeleport_Output, error) { session, _ := oatproxy.GetOAuthSession(ctx) if session == nil { return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session not found") @@ -50,7 +58,7 @@ func (s *Server) handlePlaceStreamLiveDenyTeleport(ctx context.Context, input *p return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to deny teleport") } - cancelMsg := &placestreamtypes.Livestream_TeleportCanceled{ + cancelMsg := &placestream.Livestream_TeleportCanceled{ LexiconTypeID: "place.stream.livestream#teleportCanceled", TeleportUri: input.Uri, Reason: "denied", @@ -59,7 +67,7 @@ func (s *Server) handlePlaceStreamLiveDenyTeleport(ctx context.Context, input *p s.bus.Publish(teleport.RepoDID, cancelMsg) s.bus.Publish(teleport.TargetDID, cancelMsg) - return &placestreamtypes.LiveDenyTeleport_Output{ + return &placestream.LiveDenyTeleport_Output{ Success: true, }, nil } @@ -72,7 +80,7 @@ var replicationUpgrader = websocket.Upgrader{ }, } -func (s *Server) handlePlaceStreamLiveGetSegments(ctx context.Context, before string, limit int, userDID string) (*placestreamtypes.LiveGetSegments_Output, error) { +func (s *Server) handlePlaceStreamLiveGetSegments(ctx context.Context, before string, limit int, userDID string) (*placestream.LiveGetSegments_Output, error) { if userDID == "" { return nil, echo.NewHTTPError(http.StatusBadRequest, "User DID is required") } @@ -102,7 +110,7 @@ func (s *Server) handlePlaceStreamLiveGetSegments(ctx context.Context, before st if err != nil { return nil, fmt.Errorf("error proxying to peer: %w", err) } - var output placestreamtypes.LiveGetSegments_Output + var output placestream.LiveGetSegments_Output err = json.Unmarshal(data, &output) if err != nil { return nil, fmt.Errorf("error unmarshalling response: %w", err) @@ -123,8 +131,8 @@ func (s *Server) handlePlaceStreamLiveGetSegments(ctx context.Context, before st } // Convert segments to the expected output format - output := &placestreamtypes.LiveGetSegments_Output{ - Segments: make([]*placestreamtypes.Segment_SegmentView, len(segments)), + output := &placestream.LiveGetSegments_Output{ + Segments: make([]*placestream.Segment_SegmentView, len(segments)), } for i, segment := range segments { @@ -136,9 +144,9 @@ func (s *Server) handlePlaceStreamLiveGetSegments(ctx context.Context, before st if err != nil { return nil, echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("Failed to get CID: %s", err)) } - ltd := &util.LexiconTypeDecoder{Val: record} + ltd := &lexutil.LexiconTypeDecoder{Val: record} - output.Segments[i] = &placestreamtypes.Segment_SegmentView{ + output.Segments[i] = &placestream.Segment_SegmentView{ Record: ltd, Cid: c.String(), } @@ -147,11 +155,11 @@ func (s *Server) handlePlaceStreamLiveGetSegments(ctx context.Context, before st return output, nil } -func (s *Server) handlePlaceStreamLiveGetLiveUsers(ctx context.Context, before string, limit int) (*placestreamtypes.LiveGetLiveUsers_Output, error) { +func (s *Server) handlePlaceStreamLiveGetLiveUsers(ctx context.Context, before string, limit int) (*placestream.LiveGetLiveUsers_Output, error) { // Check cache first cacheKey := fmt.Sprintf("live_users_%s_%d", before, limit) if cached, found := s.LiveUsersCache.Get(cacheKey); found { - return cached.(*placestreamtypes.LiveGetLiveUsers_Output), nil + return cached.(*placestream.LiveGetLiveUsers_Output), nil } var beforeTime *time.Time @@ -175,7 +183,7 @@ func (s *Server) handlePlaceStreamLiveGetLiveUsers(ctx context.Context, before s return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to fetch livestreams") } - streams := make([]*placestreamtypes.Livestream_LivestreamView, len(ls)) + streams := make([]*placestream.Livestream_LivestreamView, len(ls)) for i, l := range ls { stream, err := l.ToLivestreamView() @@ -183,14 +191,14 @@ func (s *Server) handlePlaceStreamLiveGetLiveUsers(ctx context.Context, before s return nil, echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("Failed to convert livestream to streamplace livestream: %s", err)) } viewers := spmetrics.GetViewCount(stream.Author.Did) - stream.ViewerCount = &placestreamtypes.Livestream_ViewerCount{ + stream.ViewerCount = &placestream.Livestream_ViewerCount{ LexiconTypeID: "place.stream.livestream#viewerCount", Count: int64(viewers), } streams[i] = stream } - liveUsers := &placestreamtypes.LiveGetLiveUsers_Output{ + liveUsers := &placestream.LiveGetLiveUsers_Output{ Streams: streams, } @@ -254,7 +262,7 @@ func (s *Server) handlePlaceStreamLiveSubscribeSegments(c echo.Context) error { } } -func (s *Server) handlePlaceStreamLiveGetRecommendations(ctx context.Context, userDID string) (*placestreamtypes.LiveGetRecommendations_Output, error) { +func (s *Server) handlePlaceStreamLiveGetRecommendations(ctx context.Context, userDID string) (*placestream.LiveGetRecommendations_Output, error) { if userDID == "" { return nil, echo.NewHTTPError(http.StatusBadRequest, "userDID is required") } @@ -275,16 +283,16 @@ func (s *Server) handlePlaceStreamLiveGetRecommendations(ctx context.Context, us } if len(liveStreamers) > 0 { - var recommendations []*placestreamtypes.LiveGetRecommendations_Output_Recommendations_Elem + var recommendations []*placestream.LiveGetRecommendations_Output_Recommendations_Elem for _, did := range liveStreamers { - recommendations = append(recommendations, &placestreamtypes.LiveGetRecommendations_Output_Recommendations_Elem{ - LiveGetRecommendations_LivestreamRecommendation: &placestreamtypes.LiveGetRecommendations_LivestreamRecommendation{ + recommendations = append(recommendations, &placestream.LiveGetRecommendations_Output_Recommendations_Elem{ + LiveGetRecommendations_LivestreamRecommendation: &placestream.LiveGetRecommendations_LivestreamRecommendation{ Did: did, Source: "streamer", }, }) } - return &placestreamtypes.LiveGetRecommendations_Output{ + return &placestream.LiveGetRecommendations_Output{ Recommendations: recommendations, UserDID: &userDID, }, nil @@ -308,16 +316,16 @@ func (s *Server) handlePlaceStreamLiveGetRecommendations(ctx context.Context, us } if len(liveFollows) > 0 { - var recommendations []*placestreamtypes.LiveGetRecommendations_Output_Recommendations_Elem + var recommendations []*placestream.LiveGetRecommendations_Output_Recommendations_Elem for _, did := range liveFollows { - recommendations = append(recommendations, &placestreamtypes.LiveGetRecommendations_Output_Recommendations_Elem{ - LiveGetRecommendations_LivestreamRecommendation: &placestreamtypes.LiveGetRecommendations_LivestreamRecommendation{ + recommendations = append(recommendations, &placestream.LiveGetRecommendations_Output_Recommendations_Elem{ + LiveGetRecommendations_LivestreamRecommendation: &placestream.LiveGetRecommendations_LivestreamRecommendation{ Did: did, Source: "follows", }, }) } - return &placestreamtypes.LiveGetRecommendations_Output{ + return &placestream.LiveGetRecommendations_Output{ Recommendations: recommendations, UserDID: &userDID, }, nil @@ -331,24 +339,276 @@ func (s *Server) handlePlaceStreamLiveGetRecommendations(ctx context.Context, us if err != nil { return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to filter default streamers") } - var recommendations []*placestreamtypes.LiveGetRecommendations_Output_Recommendations_Elem + var recommendations []*placestream.LiveGetRecommendations_Output_Recommendations_Elem for _, did := range liveDefaults { - recommendations = append(recommendations, &placestreamtypes.LiveGetRecommendations_Output_Recommendations_Elem{ - LiveGetRecommendations_LivestreamRecommendation: &placestreamtypes.LiveGetRecommendations_LivestreamRecommendation{ + recommendations = append(recommendations, &placestream.LiveGetRecommendations_Output_Recommendations_Elem{ + LiveGetRecommendations_LivestreamRecommendation: &placestream.LiveGetRecommendations_LivestreamRecommendation{ Did: did, Source: "host", }, }) } - return &placestreamtypes.LiveGetRecommendations_Output{ + return &placestream.LiveGetRecommendations_Output{ Recommendations: recommendations, UserDID: &userDID, }, nil } // No recommendations available - return &placestreamtypes.LiveGetRecommendations_Output{ - Recommendations: []*placestreamtypes.LiveGetRecommendations_Output_Recommendations_Elem{}, + return &placestream.LiveGetRecommendations_Output{ + Recommendations: []*placestream.LiveGetRecommendations_Output_Recommendations_Elem{}, UserDID: &userDID, }, nil } + +func (s *Server) handlePlaceStreamLiveStartLivestream(ctx context.Context, body *placestream.LiveStartLivestream_Input) (*placestream.LiveStartLivestream_Output, error) { + session, _ := oatproxy.GetOAuthSession(ctx) + if session != nil { + if session.DID != body.Streamer { + return nil, echo.NewHTTPError(http.StatusForbidden, "you are not the streamer") + } + } else { + svc := GetServiceAuth(ctx) + if svc != nil { + streamerSession, err := s.statefulDB.GetSessionByDID(body.Streamer) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "error getting streamer session", err) + } + if streamerSession == nil { + return nil, echo.NewHTTPError(http.StatusNotFound, "streamer session not found") + } + session = streamerSession + } else { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "you are not authorized") + } + } + + // proxy to the origin node if the streamer is broadcasting elsewhere + origin, err := s.statefulDB.GetLatestBroadcastOriginForStreamer(session.DID) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "error getting broadcast origin", err) + } + myDID := s.cli.ServerDID() + if origin != nil && origin.ServerDID != myDID { + bs, err := json.Marshal(body) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "error marshalling body", err) + } + data, err := s.ProxyServiceRequest(ctx, origin.ServerDID, "POST", "place.stream.live.startLivestream", + url.Values{}, + bytes.NewReader(bs), "application/json") + if err != nil { + return nil, err + } + var output placestream.LiveStartLivestream_Output + err = json.Unmarshal(data, &output) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "error unmarshalling response", err) + } + return &output, nil + } + + _, client := oatproxy.GetOAuthSession(ctx) + if client == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session required to start livestream") + } + + livestream := body.Livestream + now := time.Now().UTC().Format(time.RFC3339) + livestream.LexiconTypeID = "place.stream.livestream" + livestream.CreatedAt = now + livestream.LastSeenAt = &now + + if livestream.Thumb == nil { + // Step 1: get latest thumbnail from localDB and upload to user's PDS + var thumb *lexutil.LexBlob + dbThumb, err := s.localDB.LatestThumbnailForUser(session.DID) + if err != nil { + log.Error(ctx, "failed to get latest thumbnail", "err", err) + } + if dbThumb != nil { + aqt := aqtime.FromTime(dbThumb.Segment.StartTime) + fpath, err := s.cli.SegmentFilePath(session.DID, fmt.Sprintf("%s.%s", aqt.String(), dbThumb.Format)) + if err != nil { + log.Error(ctx, "failed to get thumbnail file path", "err", err) + } else { + thumbData, err := os.ReadFile(fpath) + if err != nil { + log.Error(ctx, "failed to read thumbnail file", "err", err) + } else { + mimeType := "image/jpeg" + if dbThumb.Format == "png" { + mimeType = "image/png" + } + + // Step 2: upload to user's PDS + var uploadOut comatproto.RepoUploadBlob_Output + err = client.Do(ctx, xrpc.Procedure, mimeType, "com.atproto.repo.uploadBlob", nil, bytes.NewReader(thumbData), &uploadOut) + if err != nil { + log.Error(ctx, "failed to upload thumbnail to PDS", "err", err) + } else { + thumb = uploadOut.Blob + } + } + } + } + livestream.Thumb = thumb + } + + // Step 3: create a Bluesky post announcing the livestream + repo, err := s.model.GetRepo(session.DID) + if err != nil { + log.Error(ctx, "failed to get repo", "err", err) + } + + handle := session.DID + if repo != nil && repo.Handle != "" { + handle = repo.Handle + } + + canonicalUrl := fmt.Sprintf("https://%s/%s", s.cli.BroadcasterHost, handle) + if livestream.CanonicalUrl != nil && *livestream.CanonicalUrl != "" { + canonicalUrl = *livestream.CanonicalUrl + } + + if body.CreateBlueskyPost == nil || *body.CreateBlueskyPost { + prefix := "🔴 LIVE " + suffix := " " + livestream.Title + postText := prefix + canonicalUrl + suffix + + linkStart := int64(len(prefix)) + linkEnd := linkStart + int64(len(canonicalUrl)) + + postRecord := &bsky.FeedPost{ + LexiconTypeID: "app.bsky.feed.post", + Text: postText, + CreatedAt: now, + Langs: []string{"en"}, + Facets: []*bsky.RichtextFacet{ + { + Index: &bsky.RichtextFacet_ByteSlice{ + ByteStart: linkStart, + ByteEnd: linkEnd, + }, + Features: []*bsky.RichtextFacet_Features_Elem{ + { + RichtextFacet_Link: &bsky.RichtextFacet_Link{ + LexiconTypeID: "app.bsky.richtext.facet#link", + Uri: canonicalUrl, + }, + }, + }, + }, + }, + Embed: &bsky.FeedPost_Embed{ + EmbedExternal: &bsky.EmbedExternal{ + External: &bsky.EmbedExternal_External{ + Title: fmt.Sprintf("@%s is 🔴LIVE on %s!", handle, s.cli.BroadcasterHost), + Uri: canonicalUrl, + Description: livestream.Title, + Thumb: livestream.Thumb, + }, + }, + }, + } + + postInput := comatproto.RepoCreateRecord_Input{ + Collection: "app.bsky.feed.post", + Record: &lexutil.LexiconTypeDecoder{Val: postRecord}, + Repo: session.DID, + } + var postOutput comatproto.RepoCreateRecord_Output + err = client.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.repo.createRecord", map[string]any{}, postInput, &postOutput) + if err != nil { + log.Error(ctx, "failed to create bluesky post", "err", err) + } else { + livestream.Post = &comatproto.RepoStrongRef{ + Uri: postOutput.Uri, + Cid: postOutput.Cid, + } + } + } + + // Step 4: create the place.stream.livestream record + lsInput := comatproto.RepoCreateRecord_Input{ + Collection: "place.stream.livestream", + Record: &lexutil.LexiconTypeDecoder{Val: livestream}, + Repo: session.DID, + } + var lsOutput comatproto.RepoCreateRecord_Output + err = client.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.repo.createRecord", map[string]any{}, lsInput, &lsOutput) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("failed to create livestream record: %v", err)) + } + + return &placestream.LiveStartLivestream_Output{ + Uri: lsOutput.Uri, + Cid: lsOutput.Cid, + }, nil +} + +func (s *Server) handlePlaceStreamLiveStopLivestream(ctx context.Context, body *placestream.LiveStopLivestream_Input) (*placestream.LiveStopLivestream_Output, error) { + now := time.Now().UTC().Format(util.ISO8601) + session, _ := oatproxy.GetOAuthSession(ctx) + + _, client := oatproxy.GetOAuthSession(ctx) + if client == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session required to stop livestream") + } + + livestream, err := s.model.GetLatestLivestreamForRepo(session.DID) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "error getting livestream", err) + } + + livestreamView, err := livestream.ToLivestreamView() + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "error converting livestream to view", err) + } + + livestreamRecord, ok := livestreamView.Record.Val.(*placestream.Livestream) + if !ok { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "livestream is not a streamplace livestream") + } + + if livestreamRecord.EndedAt != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, "livestream has already ended") + } + + livestreamRecord.EndedAt = &now + + aturi, err := syntax.ParseATURI(livestreamView.Uri) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "error parsing ATURI", err) + } + + var swapRecord *string + getOutput := comatproto.RepoGetRecord_Output{} + err = client.Do(ctx, xrpc.Query, "application/json", "com.atproto.repo.getRecord", map[string]any{ + "repo": session.DID, + "collection": "place.stream.livestream", + "rkey": aturi.RecordKey().String(), + }, nil, &getOutput) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "error getting livestream record", err) + } + swapRecord = getOutput.Cid + + lsInput := comatproto.RepoPutRecord_Input{ + Collection: "place.stream.livestream", + Record: &lexutil.LexiconTypeDecoder{Val: livestreamRecord}, + Rkey: aturi.RecordKey().String(), + Repo: session.DID, + SwapRecord: swapRecord, + } + var lsOutput comatproto.RepoPutRecord_Output + err = client.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.repo.putRecord", map[string]any{}, lsInput, &lsOutput) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "error updating livestream record", err) + } + + return &placestream.LiveStopLivestream_Output{ + Uri: lsOutput.Uri, + Cid: lsOutput.Cid, + }, nil +} diff --git a/pkg/spxrpc/stubs.go b/pkg/spxrpc/stubs.go index de19be2d..2b67de33 100644 --- a/pkg/spxrpc/stubs.go +++ b/pkg/spxrpc/stubs.go @@ -290,6 +290,8 @@ func (s *Server) RegisterHandlersPlaceStream(e *echo.Echo) error { 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.live.startLivestream", s.HandlePlaceStreamLiveStartLivestream) + e.POST("/xrpc/place.stream.live.stopLivestream", s.HandlePlaceStreamLiveStopLivestream) e.POST("/xrpc/place.stream.moderation.createBlock", s.HandlePlaceStreamModerationCreateBlock) e.POST("/xrpc/place.stream.moderation.createGate", s.HandlePlaceStreamModerationCreateGate) e.POST("/xrpc/place.stream.moderation.deleteBlock", s.HandlePlaceStreamModerationDeleteBlock) @@ -524,6 +526,42 @@ func (s *Server) HandlePlaceStreamLiveSearchActorsTypeahead(c echo.Context) erro return c.JSON(200, out) } +func (s *Server) HandlePlaceStreamLiveStartLivestream(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamLiveStartLivestream") + defer span.End() + + var body placestream.LiveStartLivestream_Input + if err := c.Bind(&body); err != nil { + return err + } + var out *placestream.LiveStartLivestream_Output + var handleErr error + // func (s *Server) handlePlaceStreamLiveStartLivestream(ctx context.Context,body *placestream.LiveStartLivestream_Input) (*placestream.LiveStartLivestream_Output, error) + out, handleErr = s.handlePlaceStreamLiveStartLivestream(ctx, &body) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + +func (s *Server) HandlePlaceStreamLiveStopLivestream(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamLiveStopLivestream") + defer span.End() + + var body placestream.LiveStopLivestream_Input + if err := c.Bind(&body); err != nil { + return err + } + var out *placestream.LiveStopLivestream_Output + var handleErr error + // func (s *Server) handlePlaceStreamLiveStopLivestream(ctx context.Context,body *placestream.LiveStopLivestream_Input) (*placestream.LiveStopLivestream_Output, error) + out, handleErr = s.handlePlaceStreamLiveStopLivestream(ctx, &body) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + func (s *Server) HandlePlaceStreamModerationCreateBlock(c echo.Context) error { ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamModerationCreateBlock") defer span.End() diff --git a/pkg/streamplace/livestartLivestream.go b/pkg/streamplace/livestartLivestream.go new file mode 100644 index 00000000..1d697e1d --- /dev/null +++ b/pkg/streamplace/livestartLivestream.go @@ -0,0 +1,38 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +// Lexicon schema: place.stream.live.startLivestream + +package streamplace + +import ( + "context" + + lexutil "github.com/bluesky-social/indigo/lex/util" +) + +// LiveStartLivestream_Input is the input argument to a place.stream.live.startLivestream call. +type LiveStartLivestream_Input struct { + // createBlueskyPost: Whether to create a Bluesky post announcing the livestream. + CreateBlueskyPost *bool `json:"createBlueskyPost,omitempty" cborgen:"createBlueskyPost,omitempty"` + Livestream *Livestream `json:"livestream" cborgen:"livestream"` + // streamer: The DID of the streamer. + Streamer string `json:"streamer" cborgen:"streamer"` +} + +// LiveStartLivestream_Output is the output of a place.stream.live.startLivestream call. +type LiveStartLivestream_Output struct { + // cid: The CID of the livestream record. + Cid string `json:"cid" cborgen:"cid"` + // uri: The URI of the livestream record. + Uri string `json:"uri" cborgen:"uri"` +} + +// LiveStartLivestream calls the XRPC method "place.stream.live.startLivestream". +func LiveStartLivestream(ctx context.Context, c lexutil.LexClient, input *LiveStartLivestream_Input) (*LiveStartLivestream_Output, error) { + var out LiveStartLivestream_Output + if err := c.LexDo(ctx, lexutil.Procedure, "application/json", "place.stream.live.startLivestream", nil, input, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/pkg/streamplace/livestopLivestream.go b/pkg/streamplace/livestopLivestream.go new file mode 100644 index 00000000..40d943be --- /dev/null +++ b/pkg/streamplace/livestopLivestream.go @@ -0,0 +1,33 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +// Lexicon schema: place.stream.live.stopLivestream + +package streamplace + +import ( + "context" + + lexutil "github.com/bluesky-social/indigo/lex/util" +) + +// LiveStopLivestream_Input is the input argument to a place.stream.live.stopLivestream call. +type LiveStopLivestream_Input struct { +} + +// LiveStopLivestream_Output is the output of a place.stream.live.stopLivestream call. +type LiveStopLivestream_Output struct { + // cid: The new CID of the stopped livestream record. + Cid string `json:"cid" cborgen:"cid"` + // uri: The URI of the stopped livestream record. + Uri string `json:"uri" cborgen:"uri"` +} + +// LiveStopLivestream calls the XRPC method "place.stream.live.stopLivestream". +func LiveStopLivestream(ctx context.Context, c lexutil.LexClient, input *LiveStopLivestream_Input) (*LiveStopLivestream_Output, error) { + var out LiveStopLivestream_Output + if err := c.LexDo(ctx, lexutil.Procedure, "application/json", "place.stream.live.stopLivestream", nil, input, &out); err != nil { + return nil, err + } + + return &out, nil +}