diff --git a/js/app/components/mobile/player.tsx b/js/app/components/mobile/player.tsx index 049f8387..c42559b4 100644 --- a/js/app/components/mobile/player.tsx +++ b/js/app/components/mobile/player.tsx @@ -9,6 +9,8 @@ import { PlayerUI, RotationProvider, Text, + useLivestream, + useLivestreamInfo, useLivestreamStore, usePlayerDimensions, usePlayerStore, @@ -131,6 +133,9 @@ function PlayerWithProvider( }; }, []); + const livestream = useLivestream(); + const localLivestreamURI = useLivestreamStore((x) => x.localLivestreamURI); + if (isStreamingElsewhere) { return ( @@ -190,6 +195,10 @@ function PlayerWithProvider( ); } + if (props.ingest && livestream && livestream.uri !== localLivestreamURI) { + return ; + } + const defaultHandleTeleport = (targetHandle: string, targetDID: string) => { navigation.navigate("Home", { screen: "Stream", @@ -412,3 +421,73 @@ export function PlayerInner( ); } + +export function LivestreamWarning() { + const livestream = useLivestream(); + const localLivestreamURI = useLivestreamStore((x) => x.localLivestreamURI); + const { toggleStopStream } = useLivestreamInfo(); + const navigation = useNavigation(); + const setLocalLivestreamURI = useLivestreamStore( + (x) => x.setLocalLivestreamURI, + ); + + const [loading, setLoading] = useState(false); + + if (livestream && livestream.uri !== localLivestreamURI) { + return ( + + You have an active livestream! + "{livestream.record.title}" + + + + + ); + } +} diff --git a/js/components/src/hooks/useLivestreamInfo.ts b/js/components/src/hooks/useLivestreamInfo.ts index 77dd9688..cb63b518 100644 --- a/js/components/src/hooks/useLivestreamInfo.ts +++ b/js/components/src/hooks/useLivestreamInfo.ts @@ -7,7 +7,9 @@ export function useLivestreamInfo(url?: string) { const ingest = usePlayerStore((x) => x.ingestConnectionState); const profile = useLivestreamStore((x) => x.profile); const endLivestream = useEndLivestream(); - + const setLocalLivestreamURI = useLivestreamStore( + (x) => x.setLocalLivestreamURI, + ); const createStreamRecord = useCreateStreamRecord(); const [title, setTitle] = useState(""); @@ -19,10 +21,11 @@ export function useLivestreamInfo(url?: string) { if (title !== "") { setRecordSubmitted(true); // Create the livestream record with title and custom url if available - await createStreamRecord({ + const { uri } = await createStreamRecord({ title, canonicalUrl: url || undefined, }); + setLocalLivestreamURI(uri); } } catch (error) { console.error("Error creating livestream:", error); diff --git a/js/components/src/livestream-store/livestream-state.tsx b/js/components/src/livestream-store/livestream-state.tsx index c41976b4..9f7a4ba6 100644 --- a/js/components/src/livestream-store/livestream-state.tsx +++ b/js/components/src/livestream-store/livestream-state.tsx @@ -32,6 +32,8 @@ export interface LivestreamState { setModerationPermissions: ( permissions: PlaceStreamModerationPermission.Record[], ) => void; + localLivestreamURI: string | null; + setLocalLivestreamURI: (uri: string | null) => void; } export interface LivestreamProblem { diff --git a/js/components/src/livestream-store/livestream-store.tsx b/js/components/src/livestream-store/livestream-store.tsx index eff4e30d..97e4ddb7 100644 --- a/js/components/src/livestream-store/livestream-store.tsx +++ b/js/components/src/livestream-store/livestream-store.tsx @@ -29,6 +29,8 @@ export const makeLivestreamStore = (): StoreApi => { hasReceivedSegment: false, moderationPermissions: [], setModerationPermissions: (perms) => set({ moderationPermissions: perms }), + localLivestreamURI: null, + setLocalLivestreamURI: (uri) => set({ localLivestreamURI: uri }), })); }; diff --git a/js/components/src/streamplace-store/stream.tsx b/js/components/src/streamplace-store/stream.tsx index dceb13f1..ec5881ea 100644 --- a/js/components/src/streamplace-store/stream.tsx +++ b/js/components/src/streamplace-store/stream.tsx @@ -197,7 +197,11 @@ export function useCreateStreamRecord() { createBlueskyPost: submitPost, }); - return output; + if (!output.success) { + throw new Error("Failed to start livestream"); + } + + return output.data; }; } diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index 8ff036b0..86ecd11f 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -468,7 +468,9 @@ var livestreamUpdateInterval = time.Second * 30 func (ss *StreamSession) UpdateLivestream(ctx context.Context, repoDID string) { select { case ss.livestreamUpdateChan <- struct{}{}: + log.Warn(ctx, "livestream update signal sent") default: + log.Warn(ctx, "livestream update channel full, signal already pending") // Channel full, signal already pending } } @@ -482,7 +484,7 @@ func (ss *StreamSession) livestreamUpdateLoop(ctx context.Context, repoDID strin return nil case <-ss.livestreamUpdateChan: if time.Since(ss.lastLivestreamTime) < livestreamUpdateInterval { - log.Debug(ctx, "not updating livestream, last livestream was less than 30 seconds ago") + log.Warn(ctx, "not updating livestream, last livestream was less than 30 seconds ago") continue } if err := ss.doUpdateLivestream(ctx, repoDID); err != nil { @@ -512,18 +514,6 @@ func (ss *StreamSession) doUpdateLivestream(ctx context.Context, repoDID string) if !ok { return fmt.Errorf("livestream is not a streamplace livestream") } - if lsvr.LastSeenAt == nil { - log.Debug(ctx, "livestream has no last seen at, skipping update") - return nil - } - lastSeenTime, err := time.Parse(time.RFC3339, *lsvr.LastSeenAt) - if err != nil { - return fmt.Errorf("could not parse last seen at: %w", err) - } - if time.Since(lastSeenTime) > 5*time.Minute { - log.Debug(ctx, "livestream is inactive, skipping update", "lastSeenAt", lastSeenTime) - return nil - } aturi, err := syntax.ParseATURI(lastLivestream.URI) if err != nil { diff --git a/pkg/statedb/queue_processor.go b/pkg/statedb/queue_processor.go index 858fe23c..64a6fa48 100644 --- a/pkg/statedb/queue_processor.go +++ b/pkg/statedb/queue_processor.go @@ -78,6 +78,9 @@ func (state *StatefulDB) processTask(ctx context.Context, task *AppTask) error { } func (state *StatefulDB) processFinalizeLivestreamTask(ctx context.Context, task *AppTask) error { + ctx = log.WithLogValues(ctx, "func", "processFinalizeLivestreamTask") + log.Debug(ctx, "processing finalize livestream task") + log.Warn(ctx, "processing finalize livestream task") var finalizeLivestreamTask FinalizeLivestreamTask if err := json.Unmarshal(task.Payload, &finalizeLivestreamTask); err != nil { return err diff --git a/pkg/statedb/task.go b/pkg/statedb/task.go index e971b65e..dadca32e 100644 --- a/pkg/statedb/task.go +++ b/pkg/statedb/task.go @@ -105,8 +105,8 @@ func (state *StatefulDB) DequeueTask(ctx context.Context, workerID string, taskT err := state.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error { query := tx.Where("status = ?", TaskStatusPending). Where("try_count < max_tries"). - Where("(lock_expires IS NULL OR lock_expires < ?)", time.Now()). - Where("(scheduled_at IS NULL OR scheduled_at <= ?)", time.Now()) + Where("(lock_expires IS NULL OR lock_expires < ?)", time.Now().UTC()). + Where("(scheduled_at IS NULL OR scheduled_at <= ?)", time.Now().UTC()) if len(taskTypes) > 0 { query = query.Where("type IN ?", taskTypes)