import { MusemeSourceVideo } from './overlay' import { Unlisten } from './types' /* eslint-disable no-console */ export type LiveCCSetup = { video: HTMLVideoElement canvas: HTMLCanvasElement } export type LiveCCCallbacks = { onOpen(): void onError(err: Error): void onFrame(frame: string, timestamp: number): void onMetadata(metadata: any, timestamp: number): void } export const rtcLiveSource = ( url: string, source: MusemeSourceVideo, { video, canvas }: LiveCCSetup, { onOpen, onFrame, onError, onMetadata }: LiveCCCallbacks ): Unlisten => { let timestamp: number let playheadMillis: number let playheadUtc: string let dataChannel: RTCDataChannel function setupDatachannel(channel: RTCDataChannel) { if (dataChannel) { dataChannel.close() } dataChannel = channel channel.onopen = () => { console.log('Data channel', channel, 'open!') } channel.onclose = () => { console.log('Data channel', channel, 'closed') } channel.onclosing = () => { console.log('Data channel', channel, 'closing') } channel.onmessage = (event) => { const eventJson = JSON.parse(event.data) onMetadata({ detected_objects: eventJson.data.metadata }, eventJson.time) } } const rtc = new RTCPeerConnection() rtc.addTransceiver('audio', { direction: 'recvonly' }) rtc.addTransceiver('video', { direction: 'recvonly' }) rtc.ondatachannel = (e) => { setupDatachannel(e.channel) } setupDatachannel(rtc.createDataChannel('JSON', { protocol: 'text' })) let frame: number rtc.ontrack = (e) => { const stream = e.streams[0] video.srcObject = stream const ctx = canvas.getContext('2d') const tick = () => { const nextTimestamp = playheadMillis + video.currentTime * 1000 if (nextTimestamp === timestamp) { frame = requestAnimationFrame(tick) return } timestamp = nextTimestamp canvas.width = video.width canvas.height = video.height const bbox = source.getRect() const rx = video.videoWidth / bbox.width const ry = video.videoHeight / bbox.height ctx.drawImage(video, 0, 0, video.videoWidth / rx, video.videoHeight / ry) const imageDataUrl = canvas.toDataURL('image/png') onFrame(imageDataUrl, timestamp) frame = requestAnimationFrame(tick) } tick() onOpen() } rtc .createOffer({ offerToReceiveVideo: true, offerToReceiveAudio: false, }) .then(async (offer) => { rtc.setLocalDescription(offer) try { const response = await fetch(url, { method: 'POST', headers: { 'Content-Type': 'application/sdp', }, body: offer.sdp, }) if (!response.ok) { onError(new Error('Could not fetch remote description.')) return } playheadMillis = parseInt(response.headers.get('Playhead-Millis')) playheadUtc = response.headers.get('Playhead-Utc') const answer = await response.text() if (rtc.signalingState === 'closed') { return } rtc.setRemoteDescription({ type: 'answer', sdp: answer }) } catch (err) { onError(err) } }) return () => { if (dataChannel.readyState !== 'closed') { dataChannel.close() } if (rtc.signalingState !== 'closed') { rtc.close() } cancelAnimationFrame(frame) } } export const rtcLiveSource2 = ( url: string, streamId: string, source: MusemeSourceVideo, { video, canvas }: LiveCCSetup, { onOpen, onFrame, onError, onMetadata }: LiveCCCallbacks, settings: { debug: boolean } ): Unlisten => { let dataChannel: RTCDataChannel let metadataChannel: RTCDataChannel let frame: number let latestMetadataTimestamp: number = 0 let latestPresentationTime: number = 0 let latestRtpTimestamp: number = 0 let startOffset: number = 0 const peerConnection = new RTCPeerConnection({ iceServers: [{ urls: 'stun:stun.cloudflare.com:3478' }], bundlePolicy: 'max-bundle', }) const main = async () => { const subscriberResponse = await fetch(`${url}/subscriber/create/`, { method: 'POST', headers: { 'Content-Type': 'application/json', }, body: JSON.stringify({ publisher_id: streamId, }), }) if (!subscriberResponse.ok) { try { const errorText = await subscriberResponse.text() onError(new Error(errorText)) } catch (err) { onError(err) } return } const subscriberData = await subscriberResponse.json() const subscriberSessionId = subscriberData.subscriber_session_id const publisherSessionId = subscriberData.publisher_session_id const remoteStream = new MediaStream() video.srcObject = remoteStream const frameCallback = (now, metadata) => { if (settings.debug) { console.log({ metadata, latestMetadataTimestamp, }) } latestPresentationTime = metadata.presentationTime latestRtpTimestamp = metadata.rtpTimestamp as number tick() video.requestVideoFrameCallback(frameCallback) } video.requestVideoFrameCallback(frameCallback) const tick = () => { if (startOffset === 0 && latestMetadataTimestamp > 0) { startOffset = latestRtpTimestamp - latestMetadataTimestamp } const timestamp = latestRtpTimestamp - startOffset const timestampInMillis = timestamp / 9 canvas.width = video.width canvas.height = video.height const bbox = source.getRect() const rx = video.videoWidth / bbox.width const ry = video.videoHeight / bbox.height const ctx = canvas.getContext('2d') ctx.drawImage(video, 0, 0, video.videoWidth / rx, video.videoHeight / ry) const imageDataUrl = canvas.toDataURL('image/png') onFrame(imageDataUrl, timestampInMillis) } peerConnection.ontrack = (event) => { if (settings.debug) { // eslint-disable-next-line no-console console.log('ontrack', event.track) } remoteStream.addTrack(event.track) } const pullResponse = await fetch(`${url}/subscriber/offer/`, { method: 'POST', headers: { 'Content-Type': 'application/json', }, body: JSON.stringify({ subscriber_session_id: subscriberSessionId, publisher_session_id: publisherSessionId, }), }) if (!pullResponse.ok) { const errorText = await pullResponse.text() onError(new Error(errorText)) return } const pullData = await pullResponse.json() const remoteDescription = new RTCSessionDescription( pullData.sessionDescription ) await peerConnection.setRemoteDescription(remoteDescription) const answer = await peerConnection.createAnswer() await peerConnection.setLocalDescription(answer) const renegotiateResponse = await fetch(`${url}/subscriber/renegotiate/`, { method: 'POST', headers: { 'Content-Type': 'application/json', }, body: JSON.stringify({ subscriber_session_id: subscriberSessionId, answer_sdp: answer.sdp, }), }) if (!renegotiateResponse.ok) { const errorText = await renegotiateResponse.text() onError(new Error(errorText)) return } const session = await createSession() const dataChannelResponse = await fetch( `${url}/subscriber/datachan-create/`, { method: 'POST', headers: { 'Content-Type': 'application/json', }, body: JSON.stringify({ subscriber_session_id: session.sessionId, }), } ) const dataChannelData = await dataChannelResponse.json() metadataChannel = session.peerConnection.createDataChannel( 'channel-one-subscribed', { negotiated: true, id: dataChannelData.dataChannels[0].id, } ) metadataChannel.addEventListener('message', (event) => { const data = JSON.parse(event.data) const { detected_objects, timestamp } = data latestMetadataTimestamp = timestamp const timestampInMillis = timestamp / 9 onMetadata({ detected_objects }, timestampInMillis) }) onOpen() } const createSession = async () => { const peerConnection = new RTCPeerConnection({ iceServers: [{ urls: 'stun:stun.cloudflare.com:3478' }], bundlePolicy: 'max-bundle', }) dataChannel = peerConnection.createDataChannel('server-events') const offer = await peerConnection.createOffer() await peerConnection.setLocalDescription(offer) const sessionResponse = await fetch(`${url}/session-create/`, { method: 'POST', headers: { 'Content-Type': 'application/json', }, body: JSON.stringify({ publisher_id: streamId, sdp: peerConnection.localDescription.sdp, type: peerConnection.localDescription.type, }), }) const sessionData = await sessionResponse.json() const { sessionDescription, sessionId } = sessionData const connected = new Promise((resolve, reject) => { setTimeout(reject, 5000) const iceConnectionStateChangeHandler = () => { if (settings.debug) { console.log( 'peerConnection.iceConnectionState', peerConnection.iceConnectionState ) } if (peerConnection.iceConnectionState === 'connected') { peerConnection.removeEventListener( 'iceconnectionstatechange', iceConnectionStateChangeHandler ) resolve(undefined) } } peerConnection.addEventListener( 'iceconnectionstatechange', iceConnectionStateChangeHandler ) }) await peerConnection.setRemoteDescription(sessionDescription) await connected return { peerConnection, sessionId, dataChannel, } } main() return () => { if (dataChannel && dataChannel.readyState !== 'closed') { dataChannel.close() } if (metadataChannel && metadataChannel.readyState !== 'closed') { metadataChannel.close() } if (peerConnection.signalingState !== 'closed') { peerConnection.close() } if (frame !== undefined) { cancelAnimationFrame(frame) frame = undefined } } } export const wsLiveSource = ( url: string, { onOpen, onError, onFrame, onMetadata }: LiveCCCallbacks ): Unlisten => { const socket = new WebSocket(url) socket.onopen = (event: MessageEvent) => { // console.log('onopen', event.data) onOpen() } socket.onerror = (event: ErrorEvent) => { onError(event.error) } socket.addEventListener('message', (event: MessageEvent) => { const { data: data_ } = event if (typeof data_ === 'string') { const obj = JSON.parse(data_) const { time, track, trackid, type, data } = obj if (data) { const parsedData = typeof data === 'string' ? JSON.parse(data) : data const { metadata } = parsedData if (metadata) { const parsedMetadata = JSON.parse(metadata) const { mjpg } = parsedMetadata onFrame(`data:image/jpg;base64,${mjpg}`, time) onMetadata(parsedMetadata, time) } } } }) return () => { if (socket.readyState !== socket.CLOSED) { socket.close() } } }