From c54f8311fa04b2fb2cf19ff17c1185167fc759f6 Mon Sep 17 00:00:00 2001 From: Koushik Dutta Date: Fri, 4 Feb 2022 10:28:52 -0800 Subject: [PATCH] common: webrtc to rtsp server --- common/src/wrtc-ffmpeg-source.ts | 241 +++++++++++++++++-------------- plugins/ring/package-lock.json | 4 +- plugins/ring/package.json | 2 +- plugins/ring/src/main.ts | 4 +- 4 files changed, 140 insertions(+), 111 deletions(-) diff --git a/common/src/wrtc-ffmpeg-source.ts b/common/src/wrtc-ffmpeg-source.ts index 551efb8c7..be9124d70 100644 --- a/common/src/wrtc-ffmpeg-source.ts +++ b/common/src/wrtc-ffmpeg-source.ts @@ -2,6 +2,8 @@ import { RTCAVSignalingOfferSetup, RTCAVMessage, FFMpegInput, MediaManager, Medi import { listenZeroSingleClient } from "./listen-cluster"; import { RTCPeerConnection, RTCRtpCodecParameters } from "@koush/werift"; import dgram from 'dgram'; +import { RtspServer } from "./rtsp-server"; +import { Socket } from "net"; function createSdpInput(audioPort: number, videoPort: number) { return `v=0 @@ -40,119 +42,146 @@ export function getRTCMediaStreamOptions(id: string, name: string, container: st }; } -export async function createRTCPeerConnectionSource(avsource: RTCAVSignalingOfferSetup, id: string, name: string, console: Console, sendOffer: (offer: RTCAVMessage) => Promise): Promise<{ - ffmpegInput: FFMpegInput, - peerConnection: RTCPeerConnection, -}> { - const udp = dgram.createSocket("udp4"); - const videoPort = Math.round(Math.random() * 40000 + 10000); - const audioPort = Math.round(Math.random() * 40000 + 10000); +export async function createRTCPeerConnectionSource(avsource: RTCAVSignalingOfferSetup, id: string, name: string, console: Console, sendOffer: (offer: RTCAVMessage) => Promise): Promise { + const videoPort = Math.round(Math.random() * 10000 + 30000); + const audioPort = Math.round(Math.random() * 10000 + 30000); - const sdpInput = await listenZeroSingleClient(); - sdpInput.clientPromise.then(client => { - client.write(createSdpInput(audioPort, videoPort)); - client.destroy(); - }) + const { clientPromise, port } = await listenZeroSingleClient(); - const pc = new RTCPeerConnection({ - codecs: { - audio: [ - new RTCRtpCodecParameters({ - mimeType: "audio/opus", - clockRate: 48000, - channels: 2, - }) - ], - video: [ - new RTCRtpCodecParameters({ - mimeType: "video/H264", - clockRate: 90000, - rtcpFeedback: [ - { type: "transport-cc" }, - { type: "ccm", parameter: "fir" }, - { type: "nack" }, - { type: "nack", parameter: "pli" }, - { type: "goog-remb" }, - ], - parameters: 'level-asymmetry-allowed=1;packetization-mode=1;profile-level-id=42e01f' - }) - ], - } - }); - let gotAudio = false; - let gotVideo = false; + let ai: NodeJS.Timeout; + let vi: NodeJS.Timeout; + let pc: RTCPeerConnection; + let socket: Socket; - const audioTransceiver = pc.addTransceiver("audio", avsource.audio as any); - audioTransceiver.onTrack.subscribe((track) => { - audioTransceiver.sender.replaceTrack(track); - track.onReceiveRtp.subscribe((rtp) => { - if (!gotAudio) { - gotAudio = true; - console.log('received first audio packet'); - } - udp!.send(rtp.serialize(), audioPort, "127.0.0.1"); - }); - track.onReceiveRtp.once(() => { - setInterval(() => audioTransceiver.receiver.sendRtcpPLI(track.ssrc!), 2000); - }); -}); - - const videoTransceiver = pc.addTransceiver("video", avsource.video as any); - videoTransceiver.onTrack.subscribe((track) => { - videoTransceiver.sender.replaceTrack(track); - track.onReceiveRtp.subscribe((rtp) => { - if (!gotVideo) { - gotVideo = true; - console.log('received first video packet'); - } - udp!.send(rtp.serialize(), videoPort, "127.0.0.1"); - }); - track.onReceiveRtp.once(() => { - setInterval(() => videoTransceiver.receiver.sendRtcpPLI(track.ssrc!), 2000); - }); - }); - - if (avsource.datachannel) - pc.createDataChannel(avsource.datachannel.label, avsource.datachannel.dict); - - const gatheringPromise = new Promise(resolve => pc.iceGatheringStateChange.subscribe(resolve)); - - let offer = await pc.createOffer(); - await pc.setLocalDescription(offer); - - await gatheringPromise; - - offer = await pc.createOffer(); - await pc.setLocalDescription(offer); - - const offerWithCandidates: RTCAVMessage = { - id: undefined, - candidates: [], - description: { - sdp: offer.sdp, - type: 'offer', - }, - configuration: {}, + const cleanup = () => { + pc?.close(); + socket?.destroy(); + clearInterval(ai); + clearInterval(vi); }; - console.log('offer sdp', offer.sdp); - const answer = await sendOffer(offerWithCandidates); - console.log('answer sdp', answer.description.sdp); - await pc.setRemoteDescription(answer.description as any); + clientPromise.then(async (client) => { + socket = client; + const rtspServer = new RtspServer(socket, createSdpInput(audioPort, videoPort)); + rtspServer.audioChannel = 0; + rtspServer.videoChannel = 2; + await rtspServer.handleSetup(); + const pc = new RTCPeerConnection({ + codecs: { + audio: [ + new RTCRtpCodecParameters({ + mimeType: "audio/opus", + clockRate: 48000, + channels: 2, + }) + ], + video: [ + new RTCRtpCodecParameters({ + mimeType: "video/H264", + clockRate: 90000, + rtcpFeedback: [ + { type: "transport-cc" }, + { type: "ccm", parameter: "fir" }, + { type: "nack" }, + { type: "nack", parameter: "pli" }, + { type: "goog-remb" }, + ], + parameters: 'level-asymmetry-allowed=1;packetization-mode=1;profile-level-id=42e01f' + }) + ], + } + }); + + socket.on('close', cleanup); + socket.on('error', cleanup); + pc.iceConnectionStateChange.subscribe(() => { + if (pc.iceConnectionState === 'disconnected' + || pc.iceConnectionState === 'failed' + || pc.iceConnectionState === 'closed') { + cleanup(); + } + }); + pc.connectionStateChange.subscribe(() => { + if (pc.connectionState === 'closed' + || pc.connectionState === 'disconnected' + || pc.connectionState === 'failed') { + cleanup(); + } + }) + + let gotAudio = false; + let gotVideo = false; + + const audioTransceiver = pc.addTransceiver("audio", avsource.audio as any); + audioTransceiver.onTrack.subscribe((track) => { + audioTransceiver.sender.replaceTrack(track); + track.onReceiveRtp.subscribe((rtp) => { + if (!gotAudio) { + gotAudio = true; + console.log('received first audio packet'); + } + rtspServer.sendAudio(rtp.serialize(), false); + }); + track.onReceiveRtcp.subscribe(rtcp => rtspServer.sendAudio(rtcp.serialize(), true)); + track.onReceiveRtp.once(() => ai = setInterval(() => audioTransceiver.receiver.sendRtcpPLI(track.ssrc!), 2000)); + }); + + const videoTransceiver = pc.addTransceiver("video", avsource.video as any); + videoTransceiver.onTrack.subscribe((track) => { + videoTransceiver.sender.replaceTrack(track); + track.onReceiveRtp.subscribe((rtp) => { + if (!gotVideo) { + gotVideo = true; + console.log('received first video packet'); + } + rtspServer.sendVideo(rtp.serialize(), false); + }); + track.onReceiveRtcp.subscribe(rtcp => rtspServer.sendVideo(rtcp.serialize(), true)) + track.onReceiveRtp.once(() => { + vi = setInterval(() => videoTransceiver.receiver.sendRtcpPLI(track.ssrc!), 2000); + }); + }); + + if (avsource.datachannel) + pc.createDataChannel(avsource.datachannel.label, avsource.datachannel.dict); + + const gatheringPromise = new Promise(resolve => pc.iceGatheringStateChange.subscribe(resolve)); + + let offer = await pc.createOffer(); + await pc.setLocalDescription(offer); + + await gatheringPromise; + + offer = await pc.createOffer(); + await pc.setLocalDescription(offer); + + const offerWithCandidates: RTCAVMessage = { + id: undefined, + candidates: [], + description: { + sdp: offer.sdp, + type: 'offer', + }, + configuration: {}, + }; + + console.log('offer sdp', offer.sdp); + const answer = await sendOffer(offerWithCandidates); + console.log('answer sdp', answer.description.sdp); + await pc.setRemoteDescription(answer.description as any); + }) + .catch(e => cleanup); + + const url = `rtsp://127.0.0.1:${port}`; return { - peerConnection: pc, - ffmpegInput: { - container: 'sdp', - url: sdpInput.url, - mediaStreamOptions: getRTCMediaStreamOptions(id, name, 'sdp'), - inputArguments: [ - // '-analyzeduration', '50000000', - // '-probesize', '50000000', - '-f', 'sdp', - '-i', sdpInput.url, - ] - }, + url, + mediaStreamOptions: getRTCMediaStreamOptions(id, name, 'rtsp'), + inputArguments: [ + "-rtsp_transport", "tcp", + "-max_delay", "1000000", + '-i', url, + ] }; } diff --git a/plugins/ring/package-lock.json b/plugins/ring/package-lock.json index 264c1bfee..294801585 100644 --- a/plugins/ring/package-lock.json +++ b/plugins/ring/package-lock.json @@ -1,12 +1,12 @@ { "name": "@scrypted/ring", - "version": "0.0.36", + "version": "0.0.37", "lockfileVersion": 2, "requires": true, "packages": { "": { "name": "@scrypted/ring", - "version": "0.0.36", + "version": "0.0.37", "dependencies": { "@homebridge/camera-utils": "^2.0.4", "@koush/ring-client-api": "^9.24.1", diff --git a/plugins/ring/package.json b/plugins/ring/package.json index e0d0cebf7..b07f6d8b2 100644 --- a/plugins/ring/package.json +++ b/plugins/ring/package.json @@ -37,5 +37,5 @@ "devDependencies": { "@scrypted/sdk": "file:../../sdk" }, - "version": "0.0.36" + "version": "0.0.37" } diff --git a/plugins/ring/src/main.ts b/plugins/ring/src/main.ts index 238c4413a..ba8fb4f30 100644 --- a/plugins/ring/src/main.ts +++ b/plugins/ring/src/main.ts @@ -426,12 +426,12 @@ class RingPlugin extends ScryptedDeviceBase implements BufferConverter, DevicePr break; } } - const result = await createRTCPeerConnectionSource(createRingRTCAVSignalingOfferSetup(device.signalingMime), 'default', 'MPEG-TS', device.console, async (offer) => { + const ffmpegInput = await createRTCPeerConnectionSource(createRingRTCAVSignalingOfferSetup(device.signalingMime), 'default', 'MPEG-TS', device.console, async (offer) => { const answer = await device.sendOffer(offer); device.console.log('webrtc answer', answer); return answer; }); - return Buffer.from(JSON.stringify(result.ffmpegInput)); + return Buffer.from(JSON.stringify(ffmpegInput)); } async clearTryDiscoverDevices() {