publisher.go 9.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352
  1. // The MIT License (MIT)
  2. //
  3. // # Copyright (c) 2021 Winlin
  4. //
  5. // Permission is hereby granted, free of charge, to any person obtaining a copy of
  6. // this software and associated documentation files (the "Software"), to deal in
  7. // the Software without restriction, including without limitation the rights to
  8. // use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies of
  9. // the Software, and to permit persons to whom the Software is furnished to do so,
  10. // subject to the following conditions:
  11. //
  12. // The above copyright notice and this permission notice shall be included in all
  13. // copies or substantial portions of the Software.
  14. //
  15. // THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
  16. // IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
  17. // FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
  18. // COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER
  19. // IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN
  20. // CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
  21. package janus
  22. import (
  23. "context"
  24. "fmt"
  25. "github.com/ossrs/go-oryx-lib/errors"
  26. "github.com/ossrs/go-oryx-lib/logger"
  27. "github.com/pion/interceptor"
  28. "github.com/pion/sdp/v3"
  29. "github.com/pion/webrtc/v3"
  30. "io"
  31. "net/url"
  32. "strconv"
  33. "strings"
  34. "sync"
  35. )
  36. func startPublish(ctx context.Context, r, sourceAudio, sourceVideo string, fps int, enableAudioLevel, enableTWCC bool) error {
  37. ctx = logger.WithContext(ctx)
  38. u, err := url.Parse(r)
  39. if err != nil {
  40. return errors.Wrapf(err, "Parse url %v", r)
  41. }
  42. var room int
  43. var display string
  44. if us := strings.SplitN(u.Path, "/", 3); len(us) >= 3 {
  45. if iv, err := strconv.Atoi(us[1]); err != nil {
  46. return errors.Wrapf(err, "parse %v", us[1])
  47. } else {
  48. room = iv
  49. }
  50. display = strings.Join(us[2:], "-")
  51. }
  52. logger.Tf(ctx, "Run publish url=%v, audio=%v, video=%v, fps=%v, audio-level=%v, twcc=%v",
  53. r, sourceAudio, sourceVideo, fps, enableAudioLevel, enableTWCC)
  54. // Filter for SPS/PPS marker.
  55. var aIngester *audioIngester
  56. var vIngester *videoIngester
  57. // For audio-level and sps/pps marker.
  58. // TODO: FIXME: Should share with player.
  59. webrtcNewPeerConnection := func(configuration webrtc.Configuration) (*webrtc.PeerConnection, error) {
  60. m := &webrtc.MediaEngine{}
  61. if err := m.RegisterDefaultCodecs(); err != nil {
  62. return nil, err
  63. }
  64. for _, extension := range []string{sdp.SDESMidURI, sdp.SDESRTPStreamIDURI, sdp.TransportCCURI} {
  65. if extension == sdp.TransportCCURI && !enableTWCC {
  66. continue
  67. }
  68. if err := m.RegisterHeaderExtension(webrtc.RTPHeaderExtensionCapability{URI: extension}, webrtc.RTPCodecTypeVideo); err != nil {
  69. return nil, err
  70. }
  71. }
  72. // https://github.com/pion/ion/issues/130
  73. // https://github.com/pion/ion-sfu/pull/373/files#diff-6f42c5ac6f8192dd03e5a17e9d109e90cb76b1a4a7973be6ce44a89ffd1b5d18R73
  74. for _, extension := range []string{sdp.SDESMidURI, sdp.SDESRTPStreamIDURI, sdp.AudioLevelURI} {
  75. if extension == sdp.AudioLevelURI && !enableAudioLevel {
  76. continue
  77. }
  78. if err := m.RegisterHeaderExtension(webrtc.RTPHeaderExtensionCapability{URI: extension}, webrtc.RTPCodecTypeAudio); err != nil {
  79. return nil, err
  80. }
  81. }
  82. registry := &interceptor.Registry{}
  83. if err := webrtc.RegisterDefaultInterceptors(m, registry); err != nil {
  84. return nil, err
  85. }
  86. if sourceAudio != "" {
  87. aIngester = newAudioIngester(sourceAudio)
  88. registry.Add(&rtpInteceptorFactory{aIngester.audioLevelInterceptor})
  89. }
  90. if sourceVideo != "" {
  91. vIngester = newVideoIngester(sourceVideo)
  92. registry.Add(&rtpInteceptorFactory{vIngester.markerInterceptor})
  93. }
  94. api := webrtc.NewAPI(webrtc.WithMediaEngine(m), webrtc.WithInterceptorRegistry(registry))
  95. return api.NewPeerConnection(configuration)
  96. }
  97. pc, err := webrtcNewPeerConnection(webrtc.Configuration{})
  98. if err != nil {
  99. return errors.Wrapf(err, "Create PC")
  100. }
  101. doClose := func() {
  102. if pc != nil {
  103. pc.Close()
  104. }
  105. if vIngester != nil {
  106. vIngester.Close()
  107. }
  108. if aIngester != nil {
  109. aIngester.Close()
  110. }
  111. }
  112. defer doClose()
  113. if vIngester != nil {
  114. if err := vIngester.AddTrack(pc, fps); err != nil {
  115. return errors.Wrapf(err, "Add track")
  116. }
  117. }
  118. if aIngester != nil {
  119. if err := aIngester.AddTrack(pc); err != nil {
  120. return errors.Wrapf(err, "Add track")
  121. }
  122. }
  123. offer, err := pc.CreateOffer(nil)
  124. if err != nil {
  125. return errors.Wrapf(err, "Create Offer")
  126. }
  127. if err := pc.SetLocalDescription(offer); err != nil {
  128. return errors.Wrapf(err, "Set offer %v", offer)
  129. }
  130. // Signaling API
  131. api := newJanusAPI(fmt.Sprintf("http://%v/janus", u.Host))
  132. webrtcUpCtx, webrtcUpCancel := context.WithCancel(ctx)
  133. api.onWebrtcUp = func(sender, sessionID uint64) {
  134. logger.Tf(ctx, "Event webrtcup: DTLS/SRTP done, from=(sender:%v,session:%v)", sender, sessionID)
  135. webrtcUpCancel()
  136. }
  137. api.onMedia = func(sender, sessionID uint64, mtype string, receiving bool) {
  138. logger.Tf(ctx, "Event media: %v receiving=%v, from=(sender:%v,session:%v)", mtype, receiving, sender, sessionID)
  139. }
  140. api.onSlowLink = func(sender, sessionID uint64, media string, lost uint64, uplink bool) {
  141. logger.Tf(ctx, "Event slowlink: %v lost=%v, uplink=%v, from=(sender:%v,session:%v)", media, lost, uplink, sender, sessionID)
  142. }
  143. api.onPublisher = func(sender, sessionID uint64, publishers []publisherInfo) {
  144. logger.Tf(ctx, "Event publisher: %v, from=(sender:%v,session:%v)", publishers, sender, sessionID)
  145. }
  146. api.onUnPublished = func(sender, sessionID, id uint64) {
  147. logger.Tf(ctx, "Event unpublish: %v, from=(sender:%v,session:%v)", id, sender, sessionID)
  148. }
  149. api.onLeave = func(sender, sessionID, id uint64) {
  150. logger.Tf(ctx, "Event leave: %v, from=(sender:%v,session:%v)", id, sender, sessionID)
  151. }
  152. if err := api.Create(ctx); err != nil {
  153. return errors.Wrapf(err, "create")
  154. }
  155. defer api.Close()
  156. publishHandleID, err := api.AttachPlugin(ctx)
  157. if err != nil {
  158. return errors.Wrapf(err, "attach plugin")
  159. }
  160. defer api.DetachPlugin(ctx, publishHandleID)
  161. if err := api.JoinAsPublisher(ctx, publishHandleID, room, display); err != nil {
  162. return errors.Wrapf(err, "join as publisher")
  163. }
  164. answer, err := api.Publish(ctx, publishHandleID, offer.SDP)
  165. if err != nil {
  166. return errors.Wrapf(err, "join as publisher")
  167. }
  168. defer api.UnPublish(ctx, publishHandleID)
  169. // Setup the offer-answer
  170. if err := pc.SetRemoteDescription(webrtc.SessionDescription{
  171. Type: webrtc.SDPTypeAnswer, SDP: answer,
  172. }); err != nil {
  173. return errors.Wrapf(err, "Set answer %v", answer)
  174. }
  175. logger.Tf(ctx, "State signaling=%v, ice=%v, conn=%v", pc.SignalingState(), pc.ICEConnectionState(), pc.ConnectionState())
  176. // ICE state management.
  177. pc.OnICEConnectionStateChange(func(state webrtc.ICEConnectionState) {
  178. logger.Tf(ctx, "ICE state %v", state)
  179. })
  180. pc.OnSignalingStateChange(func(state webrtc.SignalingState) {
  181. logger.Tf(ctx, "Signaling state %v", state)
  182. })
  183. if aIngester != nil {
  184. aIngester.sAudioSender.Transport().OnStateChange(func(state webrtc.DTLSTransportState) {
  185. logger.Tf(ctx, "DTLS state %v", state)
  186. })
  187. }
  188. ctx, cancel := context.WithCancel(ctx)
  189. pcDoneCtx, pcDoneCancel := context.WithCancel(context.Background())
  190. pc.OnConnectionStateChange(func(state webrtc.PeerConnectionState) {
  191. logger.Tf(ctx, "PC state %v", state)
  192. if state == webrtc.PeerConnectionStateConnected {
  193. pcDoneCancel()
  194. }
  195. if state == webrtc.PeerConnectionStateFailed || state == webrtc.PeerConnectionStateClosed {
  196. if ctx.Err() != nil {
  197. return
  198. }
  199. logger.Wf(ctx, "Close for PC state %v", state)
  200. cancel()
  201. }
  202. })
  203. // OK, DTLS/SRTP ok.
  204. select {
  205. case <-ctx.Done():
  206. return nil
  207. case <-webrtcUpCtx.Done():
  208. }
  209. // Wait for event from context or tracks.
  210. var wg sync.WaitGroup
  211. wg.Add(1)
  212. go func() {
  213. defer wg.Done()
  214. <-ctx.Done()
  215. doClose() // Interrupt the RTCP read.
  216. }()
  217. wg.Add(1)
  218. go func() {
  219. defer wg.Done()
  220. if aIngester == nil {
  221. return
  222. }
  223. select {
  224. case <-ctx.Done():
  225. case <-pcDoneCtx.Done():
  226. logger.Tf(ctx, "PC(ICE+DTLS+SRTP) done, start read audio packets")
  227. }
  228. buf := make([]byte, 1500)
  229. for ctx.Err() == nil {
  230. if _, _, err := aIngester.sAudioSender.Read(buf); err != nil {
  231. return
  232. }
  233. }
  234. }()
  235. wg.Add(1)
  236. go func() {
  237. defer wg.Done()
  238. if aIngester == nil {
  239. return
  240. }
  241. select {
  242. case <-ctx.Done():
  243. case <-pcDoneCtx.Done():
  244. logger.Tf(ctx, "PC(ICE+DTLS+SRTP) done, start ingest audio %v", sourceAudio)
  245. }
  246. // Read audio and send out.
  247. for ctx.Err() == nil {
  248. if err := aIngester.Ingest(ctx); err != nil {
  249. if errors.Cause(err) == io.EOF {
  250. logger.Tf(ctx, "EOF, restart ingest audio %v", sourceAudio)
  251. continue
  252. }
  253. logger.Wf(ctx, "Ignore audio err %+v", err)
  254. }
  255. }
  256. }()
  257. wg.Add(1)
  258. go func() {
  259. defer wg.Done()
  260. if vIngester == nil {
  261. return
  262. }
  263. select {
  264. case <-ctx.Done():
  265. case <-pcDoneCtx.Done():
  266. logger.Tf(ctx, "PC(ICE+DTLS+SRTP) done, start read video packets")
  267. }
  268. buf := make([]byte, 1500)
  269. for ctx.Err() == nil {
  270. if _, _, err := vIngester.sVideoSender.Read(buf); err != nil {
  271. return
  272. }
  273. }
  274. }()
  275. wg.Add(1)
  276. go func() {
  277. defer wg.Done()
  278. if vIngester == nil {
  279. return
  280. }
  281. select {
  282. case <-ctx.Done():
  283. case <-pcDoneCtx.Done():
  284. logger.Tf(ctx, "PC(ICE+DTLS+SRTP) done, start ingest video %v", sourceVideo)
  285. }
  286. for ctx.Err() == nil {
  287. if err := vIngester.Ingest(ctx); err != nil {
  288. if errors.Cause(err) == io.EOF {
  289. logger.Tf(ctx, "EOF, restart ingest video %v", sourceVideo)
  290. continue
  291. }
  292. logger.Wf(ctx, "Ignore video err %+v", err)
  293. }
  294. }
  295. }()
  296. wg.Wait()
  297. return nil
  298. }