diff --git a/js/app/components/mobile-player/use-webrtc.tsx b/js/app/components/mobile-player/use-webrtc.tsx new file mode 100644 index 00000000..7f66da69 --- /dev/null +++ b/js/app/components/mobile-player/use-webrtc.tsx @@ -0,0 +1,249 @@ +import { usePlayerStore, useStreamKey } from "@streamplace/components"; +import { useEffect, useRef, useState } from "react"; +import { useAppDispatch } from "store/hooks"; +import { RTCPeerConnection, RTCSessionDescription } from "./webrtc-primitives"; + +export default function useWebRTC( + endpoint: string, +): [MediaStream | null, boolean] { + const [mediaStream, setMediaStream] = useState(null); + const [stuck, setStuck] = useState(false); + + const lastChange = useRef(0); + + useEffect(() => { + const peerConnection = new RTCPeerConnection({ + bundlePolicy: "max-bundle", + }); + peerConnection.addTransceiver("video", { + direction: "recvonly", + }); + peerConnection.addTransceiver("audio", { + direction: "recvonly", + }); + peerConnection.addEventListener("track", (event) => { + const track = event.track; + if (!track) { + return; + } + setMediaStream(event.streams[0]); + }); + peerConnection.addEventListener("connectionstatechange", () => { + console.log("connection state change", peerConnection.connectionState); + if (peerConnection.connectionState === "closed") { + setStuck(true); + } + if (peerConnection.connectionState !== "connected") { + return; + } + }); + peerConnection.addEventListener("negotiationneeded", () => { + negotiateConnectionWithClientOffer(peerConnection, endpoint); + }); + + let lastFramesReceived = 0; + let lastAudioFramesReceived = 0; + + const handle = setInterval(async () => { + const stats = await peerConnection.getStats(); + stats.forEach((stat) => { + const mediaType = stat.mediaType /* web */ ?? stat.kind; /* native */ + if (stat.type === "inbound-rtp" && mediaType === "audio") { + const audioFramesReceived = stat.lastPacketReceivedTimestamp; + if (lastAudioFramesReceived !== audioFramesReceived) { + lastAudioFramesReceived = audioFramesReceived; + lastChange.current = Date.now(); + setStuck(false); + } + } + if (stat.type === "inbound-rtp" && mediaType === "video") { + const framesReceived = stat.framesReceived; + if (lastFramesReceived !== framesReceived) { + lastFramesReceived = framesReceived; + lastChange.current = Date.now(); + setStuck(false); + } + } + }); + if (Date.now() - lastChange.current > 2000) { + setStuck(true); + } + }, 200); + + return () => { + clearInterval(handle); + peerConnection.close(); + }; + }, [endpoint]); + return [mediaStream, stuck]; +} + +/** + * Performs the actual SDP exchange. + * + * 1. Constructs the client's SDP offer + * 2. Sends the SDP offer to the server, + * 3. Awaits the server's offer. + * + * SDP describes what kind of media we can send and how the server and client communicate. + * + * https://developer.mozilla.org/en-US/docs/Glossary/SDP + * https://www.ietf.org/archive/id/draft-ietf-wish-whip-01.html#name-protocol-operation + */ +export async function negotiateConnectionWithClientOffer( + peerConnection: RTCPeerConnection, + endpoint: string, + bearerToken?: string, +) { + /** https://developer.mozilla.org/en-US/docs/Web/API/RTCPeerConnection/createOffer */ + const offer = await peerConnection.createOffer({ + offerToReceiveAudio: true, + offerToReceiveVideo: true, + }); + /** https://developer.mozilla.org/en-US/docs/Web/API/RTCPeerConnection/setLocalDescription */ + await peerConnection.setLocalDescription(offer); + + /** Wait for ICE gathering to complete */ + let ofr = await waitToCompleteICEGathering(peerConnection); + if (!ofr) { + throw Error("failed to gather ICE candidates for offer"); + } + + /** + * As long as the connection is open, attempt to... + */ + while (peerConnection.connectionState !== "closed") { + try { + /** + * This response contains the server's SDP offer. + * This specifies how the client should communicate, + * and what kind of media client and server have negotiated to exchange. + */ + let response = await postSDPOffer(`${endpoint}`, ofr.sdp, bearerToken); + if (response.status === 201) { + let answerSDP = await response.text(); + if ((peerConnection.connectionState as string) === "closed") { + return; + } + await peerConnection.setRemoteDescription( + new RTCSessionDescription({ type: "answer", sdp: answerSDP }), + ); + return response.headers.get("Location"); + } else if (response.status === 405) { + console.log( + "Remember to update the URL passed into the WHIP or WHEP client", + ); + } else { + const errorMessage = await response.text(); + console.error(errorMessage); + } + } catch (e) { + console.error(`posting sdp offer failed: ${e}`); + } + + /** Limit reconnection attempts to at-most once every 5 seconds */ + await new Promise((r) => setTimeout(r, 5000)); + } +} + +async function postSDPOffer( + endpoint: string, + data: string, + bearerToken?: string, +) { + return await fetch(endpoint, { + method: "POST", + mode: "cors", + headers: { + "content-type": "application/sdp", + ...(bearerToken ? { Authorization: `Bearer ${bearerToken}` } : {}), + }, + body: data, + }); +} + +/** + * Receives an RTCPeerConnection and waits until + * the connection is initialized or a timeout passes. + * + * https://www.ietf.org/archive/id/draft-ietf-wish-whip-01.html#section-4.1 + * https://developer.mozilla.org/en-US/docs/Web/API/RTCPeerConnection/iceGatheringState + * https://developer.mozilla.org/en-US/docs/Web/API/RTCPeerConnection/icegatheringstatechange_event + */ +async function waitToCompleteICEGathering(peerConnection: RTCPeerConnection) { + return new Promise((resolve) => { + /** Wait at most 1 second for ICE gathering. */ + setTimeout(function () { + if (peerConnection.connectionState === "closed") { + return; + } + resolve(peerConnection.localDescription); + }, 1000); + peerConnection.addEventListener("icegatheringstatechange", (ev) => { + if (peerConnection.iceGatheringState === "complete") { + resolve(peerConnection.localDescription); + } + }); + }); +} + +export function useWebRTCIngest({ + endpoint, +}: { + endpoint: string; +}): [MediaStream | null, (mediaStream: MediaStream | null) => void] { + const [mediaStream, setMediaStream] = useState(null); + const ingestConnectionState = usePlayerStore((x) => x.ingestConnectionState); + const setIngestConnectionState = usePlayerStore( + (x) => x.setIngestConnectionState, + ); + const dispatch = useAppDispatch(); + const storedKey = useStreamKey(); + + const [retryTime, setRetryTime] = useState(0); + useEffect(() => { + if (!mediaStream) { + return; + } + if (!storedKey) { + return; + } + console.log("creating peer connection"); + const peerConnection = new RTCPeerConnection({ + bundlePolicy: "max-bundle", + }); + for (const track of mediaStream.getTracks()) { + console.log( + "adding track", + track.kind, + track.label, + track.enabled, + track.readyState, + ); + peerConnection.addTrack(track, mediaStream); + } + peerConnection.addEventListener("connectionstatechange", (ev) => { + setIngestConnectionState(peerConnection.connectionState); + console.log("connection state change", peerConnection.connectionState); + if (peerConnection.connectionState === "failed") { + setRetryTime(Date.now()); + } + }); + peerConnection.addEventListener("negotiationneeded", (ev) => { + negotiateConnectionWithClientOffer( + peerConnection, + endpoint, + storedKey.streamKey?.privateKey, + ); + }); + + peerConnection.addEventListener("track", (ev) => { + console.log(ev); + }); + + return () => { + peerConnection.close(); + }; + }, [endpoint, mediaStream, storedKey.streamKey?.privateKey, retryTime]); + return [mediaStream, setMediaStream]; +} diff --git a/js/app/components/mobile-player/video.native.tsx b/js/app/components/mobile-player/video.native.tsx index e4397d3d..c6cf2031 100644 --- a/js/app/components/mobile-player/video.native.tsx +++ b/js/app/components/mobile-player/video.native.tsx @@ -10,7 +10,7 @@ import { LayoutChangeEvent } from "react-native"; import { MediaStream, RTCView } from "react-native-webrtc"; import { View } from "tamagui"; import { srcToUrl } from "../player/shared"; -import useWebRTC from "../player/use-webrtc"; +import useWebRTC from "./use-webrtc"; // Add NativeIngestPlayer to the switch below! export default function VideoNative() { diff --git a/js/app/components/mobile-player/video.tsx b/js/app/components/mobile-player/video.tsx index d89a6c64..60655f32 100644 --- a/js/app/components/mobile-player/video.tsx +++ b/js/app/components/mobile-player/video.tsx @@ -11,12 +11,12 @@ import Hls from "hls.js"; import useStreamplaceNode from "hooks/useStreamplaceNode"; import { forwardRef, useCallback, useEffect, useRef, useState } from "react"; import { srcToUrl } from "../player/shared"; -import useWebRTC, { useWebRTCIngest } from "../player/use-webrtc"; import { logWebRTCDiagnostics, useWebRTCDiagnostics, } from "../player/webrtc-diagnostics"; import { checkWebRTCSupport } from "../player/webrtc-primitives"; +import useWebRTC, { useWebRTCIngest } from "./use-webrtc"; function assignVideoRef( ref: diff --git a/js/app/components/mobile-player/webrtc-primitives.tsx b/js/app/components/mobile-player/webrtc-primitives.tsx new file mode 100644 index 00000000..1074d7de --- /dev/null +++ b/js/app/components/mobile-player/webrtc-primitives.tsx @@ -0,0 +1,33 @@ +// Browser compatibility checks for WebRTC +const checkWebRTCSupport = () => { + if (typeof window === "undefined") { + throw new Error("WebRTC is not available in non-browser environments"); + } + + if (!window.RTCPeerConnection) { + throw new Error( + "RTCPeerConnection is not supported in this browser. Please use a modern browser that supports WebRTC.", + ); + } + + if (!window.RTCSessionDescription) { + throw new Error( + "RTCSessionDescription is not supported in this browser. Please use a modern browser that supports WebRTC.", + ); + } +}; + +// Check support immediately +try { + checkWebRTCSupport(); +} catch (error) { + console.error("WebRTC Compatibility Error:", error.message); +} + +export const RTCPeerConnection = window.RTCPeerConnection; +export const RTCSessionDescription = window.RTCSessionDescription; +export const WebRTCMediaStream = window.MediaStream; +export const mediaDevices = navigator.mediaDevices; + +// Export the compatibility checker for use in other components +export { checkWebRTCSupport }; diff --git a/js/components/src/livestream-store/index.tsx b/js/components/src/livestream-store/index.tsx index bc31680b..0e138296 100644 --- a/js/components/src/livestream-store/index.tsx +++ b/js/components/src/livestream-store/index.tsx @@ -1,3 +1,4 @@ export * from "./chat"; export * from "./context"; export * from "./livestream-store"; +export * from "./stream-key"; diff --git a/js/components/src/livestream-store/livestream-state.tsx b/js/components/src/livestream-store/livestream-state.tsx index f6fb33c6..ad57983a 100644 --- a/js/components/src/livestream-store/livestream-state.tsx +++ b/js/components/src/livestream-store/livestream-state.tsx @@ -15,4 +15,6 @@ export interface LivestreamState { segment: PlaceStreamSegment.Record | null; renditions: PlaceStreamDefs.Rendition[]; replyToMessage: ChatMessageViewHydrated | null; + streamKey: string | null; + setStreamKey: (key: string | null) => void; } diff --git a/js/components/src/livestream-store/livestream-store.tsx b/js/components/src/livestream-store/livestream-store.tsx index 4356b736..c42aebbb 100644 --- a/js/components/src/livestream-store/livestream-store.tsx +++ b/js/components/src/livestream-store/livestream-store.tsx @@ -16,6 +16,8 @@ export const makeLivestreamStore = (): StoreApi => { segment: null, renditions: [], replyToMessage: null, + streamKey: null, + setStreamKey: (sk) => set({ streamKey: sk }), })); }; diff --git a/js/components/src/livestream-store/stream-key.tsx b/js/components/src/livestream-store/stream-key.tsx new file mode 100644 index 00000000..58e69986 --- /dev/null +++ b/js/components/src/livestream-store/stream-key.tsx @@ -0,0 +1,95 @@ +import { bytesToMultibase, Secp256k1Keypair } from "@atproto/crypto"; +import { useEffect, useState } from "react"; +import { Platform } from "react-native"; +import { PlaceStreamKey } from "streamplace"; +import { privateKeyToAccount } from "viem/accounts"; +import { usePDSAgent } from "../streamplace-store/xrpc"; +import { useLivestreamStore } from "./livestream-store"; + +export const useStreamKey = (): { + streamKey: { + privateKey: string; + did: string; + address: string; + } | null; + error: string | null; +} => { + const pdsAgent = usePDSAgent(); + const streamKey = useLivestreamStore((state) => state.streamKey); + const setStreamKey = useLivestreamStore((state) => state.setStreamKey); + const [key, setKey] = useState(streamKey ? JSON.parse(streamKey) : null); + const [error, setError] = useState(null); + + useEffect(() => { + if (key) return; // already have key + + const generateKey = async () => { + if (!pdsAgent) { + setError("PDS Agent is not available"); + return; + } + let did = pdsAgent.did; + if (!did) { + setError("PDS Agent did is not available (not logged in?)"); + return; + } + + const keypair = await Secp256k1Keypair.create({ exportable: true }); + const exportedKey = await keypair.export(); + const didBytes = new TextEncoder().encode(did); + const combinedKey = new Uint8Array([...exportedKey, ...didBytes]); + const multibaseKey = bytesToMultibase(combinedKey, "base58btc"); + const hexKey = Array.from(exportedKey) + .map((b) => b.toString(16).padStart(2, "0")) + .join(""); + const account = privateKeyToAccount(`0x${hexKey}`); + const newKey = { + privateKey: multibaseKey, + did: keypair.did(), + address: account.address.toLowerCase(), + }; + + let platform: string = Platform.OS; + if ( + Platform.OS === "web" && + typeof window !== "undefined" && + window.navigator + ) { + let splitUA = window.navigator.userAgent + .split(" ") + .pop() + ?.split("/")[0]; + if (splitUA) { + platform = splitUA; + } + } else if (platform === "android") { + platform = "Android"; + } else if (platform === "ios") { + platform = "iOS"; + } else if (platform === "macos") { + platform = "macOS"; + } else if (platform === "windows") { + platform = "Windows"; + } + + const record: PlaceStreamKey.Record = { + signingKey: keypair.did(), + createdAt: new Date().toISOString(), + createdBy: "Streamplace on " + platform, + }; + await pdsAgent.com.atproto.repo.createRecord({ + repo: did, + collection: "place.stream.key", + record, + }); + + setStreamKey(JSON.stringify(newKey)); + setKey(newKey); + }; + + generateKey(); + // eslint-disable-next-line react-hooks/exhaustive-deps + }, [key, setStreamKey]); + + return { streamKey: key, error }; +};