84 class VideoSink :
public webrtc::VideoSinkInterface<webrtc::VideoFrame> {
86 VideoSink(webrtc::VideoTrackInterface* track) : track_(track) {
87 track_->AddOrUpdateSink(
this, webrtc::VideoSinkWants());
89 virtual ~VideoSink() { track_->RemoveSink(
this); }
92 virtual void OnFrame(
const webrtc::VideoFrame& video_frame) {
93 webrtc::scoped_refptr<webrtc::I420BufferInterface> buffer(
94 video_frame.video_frame_buffer()->ToI420());
96 buffer->height(), buffer->width());
100 webrtc::scoped_refptr<webrtc::VideoTrackInterface> track_;
103 class SetSessionDescriptionObserver
104 :
public webrtc::SetSessionDescriptionObserver {
106 static SetSessionDescriptionObserver* Create(
107 webrtc::PeerConnectionInterface* pc,
108 std::promise<const webrtc::SessionDescriptionInterface*>&
110 return new webrtc::RefCountedObject<SetSessionDescriptionObserver>(
113 virtual void OnSuccess() {
115 if (pc_->local_description()) {
116 promise_.set_value(pc_->local_description());
117 pc_->local_description()->ToString(&sdp);
118 }
else if (pc_->remote_description()) {
119 promise_.set_value(pc_->remote_description());
120 pc_->remote_description()->ToString(&sdp);
123 virtual void OnFailure(webrtc::RTCError error) {
124 utility::LogWarning(
"{}", error.message());
125 promise_.set_value(
nullptr);
129 SetSessionDescriptionObserver(
130 webrtc::PeerConnectionInterface* pc,
131 std::promise<const webrtc::SessionDescriptionInterface*>&
133 : pc_(pc), promise_(promise) {};
136 webrtc::PeerConnectionInterface* pc_;
137 std::promise<const webrtc::SessionDescriptionInterface*>& promise_;
140 class CreateSessionDescriptionObserver
141 :
public webrtc::CreateSessionDescriptionObserver {
143 static CreateSessionDescriptionObserver* Create(
144 webrtc::PeerConnectionInterface* pc,
145 std::promise<const webrtc::SessionDescriptionInterface*>&
147 return new webrtc::RefCountedObject<
148 CreateSessionDescriptionObserver>(pc, promise);
150 virtual void OnSuccess(webrtc::SessionDescriptionInterface* desc) {
152 desc->ToString(&sdp);
153 pc_->SetLocalDescription(
154 SetSessionDescriptionObserver::Create(pc_, promise_), desc);
156 virtual void OnFailure(webrtc::RTCError error) {
157 utility::LogWarning(
"{}", error.message());
158 promise_.set_value(
nullptr);
162 CreateSessionDescriptionObserver(
163 webrtc::PeerConnectionInterface* pc,
164 std::promise<const webrtc::SessionDescriptionInterface*>&
166 : pc_(pc), promise_(promise) {};
169 webrtc::PeerConnectionInterface* pc_;
170 std::promise<const webrtc::SessionDescriptionInterface*>& promise_;
173 class PeerConnectionStatsCollectorCallback
174 :
public webrtc::RTCStatsCollectorCallback {
176 PeerConnectionStatsCollectorCallback() {}
177 void clearReport() { report_.clear(); }
178 Json::Value getReport() {
return report_; }
181 virtual void OnStatsDelivered(
182 const webrtc::scoped_refptr<const webrtc::RTCStatsReport>&
184 for (
const webrtc::RTCStats& stats : *report) {
185 Json::Value stats_members;
186 for (
const webrtc::Attribute& attribute : stats.Attributes()) {
187 stats_members[attribute.name()] = attribute.ToString();
189 report_[stats.id()] = stats_members;
196 class DataChannelObserver :
public webrtc::DataChannelObserver {
199 webrtc::scoped_refptr<webrtc::DataChannelInterface>
201 const std::string& peerid)
202 : peer_connection_manager_(peer_connection_manager),
203 data_channel_(data_channel),
205 data_channel_->RegisterObserver(
this);
207 virtual ~DataChannelObserver() { data_channel_->UnregisterObserver(); }
210 virtual void OnStateChange() {
212 const std::string label = data_channel_->label();
213 const std::string state =
214 webrtc::DataChannelInterface::DataStateString(
215 data_channel_->state());
217 "DataChannelObserver::OnStateChange label: {}, state: {}, "
219 label, state, peerid_);
220 std::string msg(label +
" " + state);
221 webrtc::DataBuffer buffer(msg);
222 data_channel_->Send(buffer);
227 if (label ==
"ClientDataChannel" && state ==
"open") {
229 std::lock_guard<std::mutex> mutex_lock(
230 peer_connection_manager_
232 peer_connection_manager_->peerid_data_channel_ready_.insert(
235 peer_connection_manager_->SendInitFramesToPeer(peerid_);
237 if (label ==
"ClientDataChannel" &&
238 (state ==
"closed" || state ==
"closing")) {
239 std::lock_guard<std::mutex> mutex_lock(
240 peer_connection_manager_->peerid_data_channel_mutex_);
241 peer_connection_manager_->peerid_data_channel_ready_.erase(
245 virtual void OnMessage(
const webrtc::DataBuffer& buffer) {
246 std::string msg((
const char*)buffer.data.data(),
248 utility::LogDebug(
"DataChannelObserver::OnMessage: {}, msg: {}.",
249 data_channel_->label(), msg);
253 if (!reply.empty()) {
254 webrtc::DataBuffer buffer(reply);
255 data_channel_->Send(buffer);
261 webrtc::scoped_refptr<webrtc::DataChannelInterface> data_channel_;
262 const std::string peerid_;
265 class PeerConnectionObserver :
public webrtc::PeerConnectionObserver {
268 const std::string& peerid);
270 void Initialize(webrtc::scoped_refptr<webrtc::PeerConnectionInterface>
273 virtual ~PeerConnectionObserver() {
274 delete local_channel_;
275 delete remote_channel_;
283 Json::Value GetIceCandidateList() {
return ice_candidate_list_; }
285 Json::Value GetStats() {
286 stats_callback_->clearReport();
287 pc_->GetStats(stats_callback_.get());
289 while ((stats_callback_->getReport().empty()) && (--
count > 0)) {
290 std::this_thread::sleep_for(std::chrono::milliseconds(1000));
292 return Json::Value(stats_callback_->getReport());
295 webrtc::scoped_refptr<webrtc::PeerConnectionInterface>
296 GetPeerConnection() {
301 virtual void OnAddStream(
302 webrtc::scoped_refptr<webrtc::MediaStreamInterface> stream) {
303 utility::LogDebug(
"[{}] GetVideoTracks().size(): {}.",
305 webrtc::VideoTrackVector videoTracks = stream->GetVideoTracks();
306 if (videoTracks.size() > 0) {
307 video_sink_.reset(
new VideoSink(videoTracks.at(0).get()));
310 virtual void OnRemoveStream(
311 webrtc::scoped_refptr<webrtc::MediaStreamInterface> stream) {
314 virtual void OnDataChannel(
315 webrtc::scoped_refptr<webrtc::DataChannelInterface> channel) {
317 "PeerConnectionObserver::OnDataChannel peerid: {}",
319 remote_channel_ =
new DataChannelObserver(peer_connection_manager_,
322 virtual void OnRenegotiationNeeded() {
323 std::lock_guard<std::mutex> mutex_lock(
324 peer_connection_manager_->peerid_data_channel_mutex_);
325 peer_connection_manager_->peerid_data_channel_ready_.erase(peerid_);
327 "PeerConnectionObserver::OnRenegotiationNeeded peerid: {}",
330 virtual void OnIceCandidate(
331 const webrtc::IceCandidateInterface* candidate);
333 virtual void OnSignalingChange(
334 webrtc::PeerConnectionInterface::SignalingState state) {
335 utility::LogDebug(
"state: {}, peerid: {}", state, peerid_);
337 virtual void OnIceConnectionChange(
338 webrtc::PeerConnectionInterface::IceConnectionState state) {
340 webrtc::PeerConnectionInterface::kIceConnectionFailed) ||
342 webrtc::PeerConnectionInterface::kIceConnectionClosed)) {
343 ice_candidate_list_.clear();
345 std::thread([
this]() {
346 peer_connection_manager_->HangUp(peerid_);
352 virtual void OnIceGatheringChange(
353 webrtc::PeerConnectionInterface::IceGatheringState) {}
357 const std::string peerid_;
358 webrtc::scoped_refptr<webrtc::PeerConnectionInterface> pc_;
359 DataChannelObserver* local_channel_;
360 DataChannelObserver* remote_channel_;
361 Json::Value ice_candidate_list_;
362 webrtc::scoped_refptr<PeerConnectionStatsCollectorCallback>
364 std::unique_ptr<VideoSink> video_sink_;
370 const Json::Value& config,
371 const std::string& publish_filter,
372 const std::string& webrtc_udp_port_range);
376 const std::map<std::string, HttpServerRequestHandler::HttpFunction>
381 const Json::Value& json_message);
383 const Json::Value
HangUp(
const std::string& peerid);
384 const Json::Value
Call(
const std::string& peerid,
385 const std::string& window_uid,
386 const std::string& options,
387 const Json::Value& json_message);
394 void OnFrame(
const std::string& window_uid,
395 const std::shared_ptr<core::Tensor>& im);
399 const std::string& window_uid);
401 bool AddStreams(webrtc::PeerConnectionInterface* peer_connection,
402 const std::string& window_uid,
403 const std::string& options);
405 const std::string& window_uid,
406 const std::map<std::string, std::string>& opts);
409 const std::string& peerid);
413 webrtc::scoped_refptr<webrtc::PeerConnectionFactoryInterface>
417 std::unordered_map<std::string, PeerConnectionObserver*>
425 std::unordered_map<std::string,
426 webrtc::scoped_refptr<BitmapTrackSourceInterface>>
431 std::unordered_map<std::string, std::set<std::string>>
440 std::map<std::string, HttpServerRequestHandler::HttpFunction>
func_;
447 std::unordered_map<std::string, std::shared_ptr<core::Tensor>>