ipc_media.go 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453
  1. package main
  2. import (
  3. "errors"
  4. "net/http"
  5. "os"
  6. "strconv"
  7. "strings"
  8. "time"
  9. "github.com/go-resty/resty/v2"
  10. "github.com/sirupsen/logrus"
  11. "lc/common/mqtt"
  12. "lc/common/onvif/profiles/media"
  13. "lc/common/onvif/soap"
  14. "lc/common/protocol"
  15. "lc/common/util"
  16. )
  17. var (
  18. APPLIVE = "live"
  19. UploadSnapshot = "/camera/v1/snapshot"
  20. UploadFlv = "/camera/v1/flv"
  21. UploadPresetSnapshot = "/camera/v1/preset/picture"
  22. )
  23. func (o *LcDevice) GetProfiles() error {
  24. XAddr, ok := o.endpoints["media"]
  25. if !ok {
  26. return errors.New("不支持media")
  27. }
  28. m := media.NewMedia(o.client, XAddr)
  29. reply, err := m.GetProfiles(&media.GetProfiles{})
  30. if err != nil {
  31. if serr, ok := err.(*soap.SOAPFault); ok {
  32. logrus.Error(serr)
  33. }
  34. return err
  35. }
  36. for _, v := range reply.Profiles {
  37. vt := VideoToken{
  38. Token: string(v.Token),
  39. Name: string(v.Name),
  40. VideoEncoding: string(v.VideoEncoderConfiguration.Encoding),
  41. VideoResolution: Resolution{
  42. Width: v.VideoEncoderConfiguration.Resolution.Width,
  43. Height: v.VideoEncoderConfiguration.Resolution.Height,
  44. },
  45. VideoFrameRateLimit: v.VideoEncoderConfiguration.RateControl.FrameRateLimit,
  46. VideoBitrateLimit: v.VideoEncoderConfiguration.RateControl.BitrateLimit,
  47. }
  48. o.tokens[vt.Token] = &vt
  49. }
  50. return nil
  51. }
  52. func (o *LcDevice) GetAllStreamUri() error {
  53. XAddr, ok := o.endpoints["media"]
  54. if !ok {
  55. return errors.New("不支持media")
  56. }
  57. m := media.NewMedia(o.client, XAddr)
  58. for k, v := range o.tokens {
  59. for retry := 3; retry > 0; retry-- {
  60. reply, err := m.GetStreamUri(&media.GetStreamUri{ProfileToken: media.ReferenceToken(k)})
  61. if err != nil {
  62. //soap错误,则退出
  63. if serr, ok := err.(*soap.SOAPFault); ok {
  64. logrus.Error(serr)
  65. break
  66. }
  67. time.Sleep(3 * time.Second)
  68. } else {
  69. v.VideoRtspUrl = strings.ReplaceAll(string(reply.MediaUri.Uri), "rtsp://", "rtsp://"+o.onvifDev.User+":"+o.onvifDev.Password+"@")
  70. break
  71. }
  72. }
  73. }
  74. return nil
  75. }
  76. func (o *LcDevice) GetAllSnapshotUri() error {
  77. XAddr, ok := o.endpoints["media"]
  78. if !ok {
  79. return errors.New("不支持media")
  80. }
  81. m := media.NewMedia(o.client, XAddr)
  82. for k, v := range o.tokens {
  83. for retry := 3; retry > 0; retry-- {
  84. reply, err := m.GetSnapshotUri(&media.GetSnapshotUri{ProfileToken: media.ReferenceToken(k)})
  85. if err != nil {
  86. //soap错误,则退出
  87. if serr, ok := err.(*soap.SOAPFault); ok {
  88. logrus.Error(serr)
  89. break
  90. }
  91. //其他错误停3秒重试
  92. time.Sleep(5 * time.Second)
  93. } else {
  94. v.SnapshotUrl = string(reply.MediaUri.Uri)
  95. break
  96. }
  97. }
  98. }
  99. return nil
  100. }
  101. // SnapshotFromRtsp 通过RTSP截取一张图片
  102. func (o *LcDevice) SnapshotFromRtsp(token string) (string, error) {
  103. RtspUrl := o.GetRtspUrl(token)
  104. if RtspUrl == "" {
  105. logrus.Errorf("未获取到设备[%s]的RtspUrl,无法截图.", o.onvifDev.Code)
  106. return "", errors.New("找不到RtspUrl")
  107. }
  108. wdir, _ := os.Getwd()
  109. proc := &os.ProcAttr{
  110. Dir: wdir,
  111. Env: os.Environ(),
  112. Files: []*os.File{
  113. os.Stdin,
  114. os.Stdout,
  115. os.Stderr,
  116. },
  117. }
  118. fileName := util.GetPath(4) + o.onvifDev.Code + "_" + strconv.FormatInt(util.MlNow().Unix(), 10) + ".jpg"
  119. args := []string{onvifDevConfig.Ffmpeg, "-i", RtspUrl, "-t", "0.001", "-y", "-f", "mjpeg", "-r", "1", fileName}
  120. process, err := os.StartProcess(onvifDevConfig.Ffmpeg, args, proc)
  121. if err != nil {
  122. logrus.Errorf("启动ffmpeg失败:%s", err.Error())
  123. return "", err
  124. } else {
  125. go func() {
  126. time.Sleep(6 * time.Second)
  127. if process != nil {
  128. process.Signal(os.Kill)
  129. }
  130. }()
  131. _, err = process.Wait()
  132. process = nil
  133. }
  134. //判断是否已经抓图成功
  135. if FileExist(fileName) {
  136. return fileName, nil
  137. }
  138. err = errors.New("snapshotFromRtsp:抓图失败")
  139. return "", err
  140. }
  141. func (o *LcDevice) GetRtspUrl(token string) string {
  142. rtspurl := ""
  143. if v, ok := o.tokens[token]; !ok {
  144. max := int32(0)
  145. for _, v := range o.tokens { //查找分辨率最高的SnapshotUrl
  146. r := v.VideoResolution.Height * v.VideoResolution.Width
  147. if r > max {
  148. max = r
  149. rtspurl = v.VideoRtspUrl
  150. }
  151. }
  152. } else {
  153. rtspurl = v.VideoRtspUrl
  154. }
  155. return rtspurl
  156. }
  157. func (o *LcDevice) GetSnapshotUrl(token string) string {
  158. snapshoturl := ""
  159. max := int32(0)
  160. if v, ok := o.tokens[token]; !ok {
  161. for _, v := range o.tokens { //查找分辨率最高的SnapshotUrl
  162. r := v.VideoResolution.Height * v.VideoResolution.Width
  163. if r > max {
  164. max = r
  165. snapshoturl = v.SnapshotUrl
  166. }
  167. }
  168. } else {
  169. snapshoturl = v.SnapshotUrl
  170. }
  171. return snapshoturl
  172. }
  173. func (o *LcDevice) GetProfileToken(token string) string {
  174. max := int32(0)
  175. if _, ok := o.tokens[token]; !ok {
  176. for _, v := range o.tokens { //查找分辨率最高的SnapshotUrl
  177. r := v.VideoResolution.Height * v.VideoResolution.Width
  178. if r > max {
  179. max = r
  180. token = v.Token
  181. }
  182. }
  183. }
  184. return token
  185. }
  186. func (o *LcDevice) Snapshot(token string) {
  187. snapshoturl := o.GetSnapshotUrl(token)
  188. if snapshoturl == "" {
  189. logrus.Debugf("未获取到设备[%s]的SnapshotUrl,无法截图.", o.onvifDev.Code)
  190. return
  191. }
  192. var fileName string
  193. client := resty.New()
  194. client.SetTimeout(6 * time.Second)
  195. client.SetBasicAuth(o.onvifDev.User, o.onvifDev.Password)
  196. resp, err := client.R().Get(snapshoturl)
  197. if err != nil {
  198. logrus.Errorf("抓图失败:%s", err.Error())
  199. fileName, err = o.SnapshotFromRtsp(token)
  200. } else {
  201. if resp.StatusCode() != http.StatusOK {
  202. logrus.Errorf("抓图失败,http返回错误:%s", resp.Status())
  203. fileName, err = o.SnapshotFromRtsp(token)
  204. } else if len(resp.Body()) == 0 {
  205. logrus.Error("抓图失败,内容为空")
  206. fileName, err = o.SnapshotFromRtsp(token)
  207. } else {
  208. fileName = util.GetPath(4) + o.onvifDev.Code + "_" + strconv.FormatInt(util.MlNow().Unix(), 10) + ".jpg"
  209. if err = os.WriteFile(fileName, resp.Body(), os.ModePerm); err != nil {
  210. logrus.Errorf("存储图片失败:%s", err.Error())
  211. return
  212. }
  213. }
  214. }
  215. if fileName != "" && err == nil {
  216. resp, err := client.R().SetFile("file", fileName).Post(o.onvifDev.WebServer + UploadSnapshot)
  217. if err != nil {
  218. logrus.Errorf("上传文件失败:%s", err.Error())
  219. } else {
  220. if resp.StatusCode() != http.StatusOK {
  221. logrus.Errorf("上传文件失败,http返回错误:%s", resp.Status())
  222. }
  223. }
  224. os.Remove(fileName)
  225. }
  226. }
  227. func (o *LcDevice) Snapshot2(file string) error {
  228. snapshoturl := o.GetSnapshotUrl("")
  229. if snapshoturl == "" {
  230. return errors.New("找不到链接,截图错误")
  231. }
  232. var fileName string
  233. client := resty.New()
  234. client.SetTimeout(6 * time.Second)
  235. client.SetBasicAuth(o.onvifDev.User, o.onvifDev.Password)
  236. resp, err := client.R().Get(snapshoturl)
  237. if err != nil {
  238. return err
  239. }
  240. if resp.StatusCode() != http.StatusOK {
  241. return errors.New("调用截图接口发生错误")
  242. }
  243. fileName = util.GetPath(4) + file
  244. err = os.WriteFile(fileName, resp.Body(), os.ModePerm)
  245. if err != nil {
  246. return err
  247. }
  248. resp, err = client.R().SetFile("file", fileName).Post(o.onvifDev.WebServer + UploadPresetSnapshot)
  249. os.Remove(fileName)
  250. if err != nil {
  251. return err
  252. }
  253. if resp.StatusCode() != http.StatusOK {
  254. return errors.New("截图上传发生错误")
  255. }
  256. return nil
  257. }
  258. func (o *LcDevice) CreateProcess(rtspurl string) {
  259. wdir, _ := os.Getwd()
  260. proc := &os.ProcAttr{
  261. Dir: wdir,
  262. Env: os.Environ(),
  263. Files: []*os.File{
  264. os.Stdin,
  265. os.Stdout,
  266. os.Stderr,
  267. },
  268. }
  269. args := []string{onvifDevConfig.Ffmpeg, "-fflags", "nobuffer", "-i", rtspurl, "-c", "copy", "-f", "flv", o.onvifDev.RtmpServer + "/" + APPLIVE + "/" + o.onvifDev.Code}
  270. process, err := os.StartProcess(onvifDevConfig.Ffmpeg, args, proc)
  271. if err != nil {
  272. logrus.Errorf("启动ffmpeg失败:%s", err.Error())
  273. return
  274. }
  275. o.ffmpeg = process
  276. go o.WatchFfmpeg()
  277. logrus.Infof("启动ffmpeg成功:pid=%d", o.ffmpeg.Pid)
  278. }
  279. func (o *LcDevice) CleanupProcess() {
  280. if o.ffmpeg != nil {
  281. logrus.Infof("ffmpeg进程销毁:pid=%d", o.ffmpeg.Pid)
  282. if err := o.ffmpeg.Kill(); err != nil {
  283. logrus.Infof("ffmpeg进程%dKill失败:%s", o.ffmpeg.Pid, err.Error())
  284. }
  285. }
  286. }
  287. func (o *LcDevice) WatchFfmpeg() {
  288. status := make(chan *os.ProcessState)
  289. died := make(chan error)
  290. go func() {
  291. state, err := o.ffmpeg.Wait()
  292. if err != nil {
  293. died <- err
  294. return
  295. }
  296. status <- state
  297. }()
  298. select {
  299. case s := <-status:
  300. logrus.Infof("ffmpeg已退出:%s", s.String())
  301. case err := <-died:
  302. logrus.Infof("ffmpeg已退出:%s", err.Error())
  303. }
  304. o.ffmpeg = nil
  305. }
  306. func (o *LcDevice) RecordToFLV(token string, second int) {
  307. //最长录制60秒视频
  308. if second > 60 {
  309. second = 60
  310. }
  311. RtspUrl := o.GetRtspUrl(token)
  312. if RtspUrl == "" {
  313. logrus.Debugf("未获取到设备[%s]的RtspUrl,无法录制视频.", o.onvifDev.Code)
  314. return
  315. }
  316. wdir, _ := os.Getwd()
  317. proc := &os.ProcAttr{
  318. Dir: wdir,
  319. Env: os.Environ(),
  320. Files: []*os.File{
  321. os.Stdin,
  322. os.Stdout,
  323. os.Stderr,
  324. },
  325. }
  326. fileName := util.GetPath(4) + o.onvifDev.Code + "_" + strconv.FormatInt(util.MlNow().Unix(), 10) + ".flv"
  327. args := []string{onvifDevConfig.Ffmpeg, "-i", RtspUrl, "-c", "copy", "-t", strconv.Itoa(second), fileName}
  328. process, err := os.StartProcess(onvifDevConfig.Ffmpeg, args, proc)
  329. if err != nil {
  330. logrus.Errorf("启动ffmpeg失败:%s", err.Error())
  331. return
  332. } else if process != nil {
  333. go func() {
  334. time.Sleep(time.Duration(second+5) * time.Second)
  335. if process != nil {
  336. process.Signal(os.Kill)
  337. }
  338. }()
  339. _, err = process.Wait()
  340. //判断是否已经抓图成功
  341. if fileName != "" && err == nil {
  342. if FileExist(fileName) {
  343. client := resty.New()
  344. client.SetTimeout(6 * time.Second)
  345. resp, err := client.R().SetFile("file", fileName).Post(o.onvifDev.WebServer + UploadFlv)
  346. if err != nil {
  347. logrus.Errorf("上传文件失败:%s", err.Error())
  348. } else {
  349. if resp.StatusCode() != http.StatusOK {
  350. logrus.Errorf("上传文件失败,http返回错误:%s", resp.Status())
  351. }
  352. }
  353. os.Remove(fileName)
  354. }
  355. }
  356. }
  357. }
  358. // HandleTpOnvifSnapshot 截图
  359. func (o *LcDevice) HandleTpOnvifSnapshot(m mqtt.Message) {
  360. var obj protocol.Pack_MediaCommonInfo
  361. if err := obj.DeCode(m.PayloadString()); err != nil {
  362. return
  363. }
  364. if o.onvifDev.Code != obj.Id {
  365. return
  366. }
  367. var ret protocol.Pack_IPCCommonACK
  368. if strRet, err := ret.EnCode(o.onvifDev.Code, appConfig.GID, "", obj.Seq, nil); err == nil {
  369. GetMQTTMgr().Publish(GetTopic(o.GetDevType(), o.onvifDev.Code, protocol.TP_ONVIF_SNAPSHOT_ACK), strRet, mqtt.AtMostOnce, ToAll)
  370. }
  371. go o.Snapshot(obj.Data.ProfileToken)
  372. }
  373. // HandleTpOnvifVideo 拉流推流,停止拉流推流
  374. func (o *LcDevice) HandleTpOnvifVideo(m mqtt.Message) {
  375. var obj protocol.Pack_MediaCommonInfo
  376. if err := obj.DeCode(m.PayloadString()); err != nil {
  377. return
  378. }
  379. if o.onvifDev.Code != obj.Id {
  380. return
  381. }
  382. var ret protocol.Pack_IPCCommonACK
  383. if strRet, err := ret.EnCode(o.onvifDev.Code, appConfig.GID, "", obj.Seq, nil); err == nil {
  384. GetMQTTMgr().Publish(GetTopic(o.GetDevType(), o.onvifDev.Code, protocol.TP_ONVIF_VIDEO_ACK), strRet, mqtt.AtMostOnce, ToAll)
  385. }
  386. //如果已开启GB28181,则不用启动ffmpeg推流
  387. if o.onvifDev.Gb28181 {
  388. logrus.Debugf("设备[%s]已开启GB28181视频服务,无需ffmpeg推流", o.onvifDev.Code)
  389. return
  390. }
  391. if obj.Data.Flag == 1 {
  392. o.observer++
  393. logrus.Debugf("执行推流一次,当前有%d次引用", o.observer)
  394. if o.ffmpeg != nil {
  395. return
  396. }
  397. if rtspurl := o.GetRtspUrl(obj.Data.ProfileToken); rtspurl == "" {
  398. logrus.Errorf("请求设备[%s]的视频流错误,找不到拉流地址.", o.onvifDev.Code)
  399. return
  400. } else {
  401. o.CreateProcess(rtspurl)
  402. }
  403. } else if obj.Data.Flag == 2 { //停止推流
  404. if o.observer > 0 {
  405. o.observer--
  406. logrus.Debugf("停止推流一次,余下%d次引用", o.observer)
  407. if o.observer == 0 {
  408. o.CleanupProcess()
  409. }
  410. }
  411. }
  412. }
  413. func (o *LcDevice) HandleTpOnvifRecord(m mqtt.Message) {
  414. var obj protocol.Pack_MediaCommonInfo
  415. if err := obj.DeCode(m.PayloadString()); err != nil {
  416. return
  417. }
  418. if o.onvifDev.Code != obj.Id {
  419. return
  420. }
  421. var ret protocol.Pack_IPCCommonACK
  422. if strRet, err := ret.EnCode(o.onvifDev.Code, appConfig.GID, "", obj.Seq, nil); err == nil {
  423. GetMQTTMgr().Publish(GetTopic(o.GetDevType(), o.onvifDev.Code, protocol.TP_ONVIF_RECORD_ACK), strRet, mqtt.AtMostOnce, ToAll)
  424. }
  425. go o.RecordToFLV("", 30)
  426. }