process.go 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369
  1. package service
  2. import (
  3. "bytes"
  4. "context"
  5. "encoding/json"
  6. "errors"
  7. "fmt"
  8. "github.com/fsnotify/fsnotify"
  9. "gorm.io/gorm"
  10. "io"
  11. "log"
  12. "mime/multipart"
  13. "net/http"
  14. "os"
  15. "os/exec"
  16. "path/filepath"
  17. "speechAnalysis/conf"
  18. "speechAnalysis/constvar"
  19. "speechAnalysis/models"
  20. "speechAnalysis/pkg/logx"
  21. "strings"
  22. "time"
  23. )
  24. // Response 结构体用于存储响应体的内容
  25. type Response struct {
  26. Code int `json:"code"`
  27. Msg string `json:"msg"`
  28. Result string `json:"result"`
  29. Score float64 `json:"score"`
  30. }
  31. func AnalysisAudio(filename string, targetURL string) (resp Response, err error) {
  32. file, err := os.Open(filename)
  33. if err != nil {
  34. return
  35. }
  36. defer file.Close()
  37. // 创建一个缓冲区来存储表单数据
  38. var requestBody bytes.Buffer
  39. writer := multipart.NewWriter(&requestBody)
  40. // 创建一个表单字段,用于存储文件
  41. fileWriter, err := writer.CreateFormFile("audio", filename)
  42. if err != nil {
  43. return
  44. }
  45. // 将文件内容复制到表单字段中
  46. _, err = io.Copy(fileWriter, file)
  47. if err != nil {
  48. return
  49. }
  50. // 关闭表单写入器,以便写入末尾的边界
  51. writer.Close()
  52. // 创建POST请求,指定URL和请求体
  53. request, err := http.NewRequest("POST", targetURL, &requestBody)
  54. if err != nil {
  55. return
  56. }
  57. // 设置请求头,指定Content-Type为multipart/form-data
  58. request.Header.Set("Content-Type", writer.FormDataContentType())
  59. // 发送请求
  60. client := &http.Client{}
  61. response, err := client.Do(request)
  62. if err != nil {
  63. return
  64. }
  65. defer response.Body.Close()
  66. // 读取响应
  67. body := &bytes.Buffer{}
  68. _, err = io.Copy(body, response.Body)
  69. if err != nil {
  70. return
  71. }
  72. err = json.NewDecoder(body).Decode(&resp)
  73. if err != nil {
  74. return
  75. }
  76. return
  77. }
  78. func Process(audioId uint) (err error) {
  79. audio, err := models.NewAudioSearch().SetID(audioId).First()
  80. if err != nil {
  81. return errors.New("查找音频失败")
  82. }
  83. if audio.AudioStatus != constvar.AudioStatusUploadOk && audio.AudioStatus != constvar.AudioStatusFailed {
  84. return errors.New("状态不正确")
  85. }
  86. err = models.NewAudioSearch().SetID(audioId).UpdateByMap(map[string]interface{}{"audio_status": constvar.AudioStatusProcessing})
  87. if err != nil {
  88. return errors.New("DB错误")
  89. }
  90. go func() {
  91. var resp Response
  92. sz := audio.Size / 1024 / 1024
  93. if sz > 2 {
  94. resp, err = AnalysisAudio(audio.FilePath, conf.AanlysisConf.LongUrl)
  95. } else {
  96. resp, err = AnalysisAudio(audio.FilePath, conf.AanlysisConf.Url)
  97. }
  98. if err != nil {
  99. logx.Errorf("err when AnalysisAudio:%v", err)
  100. _ = models.NewAudioSearch().SetID(audioId).UpdateByMap(map[string]interface{}{"audio_status": constvar.AudioStatusFailed})
  101. return
  102. }
  103. if resp.Code != 0 {
  104. logx.Errorf("AnalysisAudio error return:%v", resp)
  105. _ = models.NewAudioSearch().SetID(audioId).UpdateByMap(map[string]interface{}{"audio_status": constvar.AudioStatusFailed})
  106. return
  107. }
  108. logx.Infof("AnalysisAudio result: %v", resp)
  109. words := GetWordFromText(resp.Result, audio)
  110. err = models.WithTransaction(func(db *gorm.DB) error {
  111. err = models.NewAudioSearch().SetOrm(db).SetID(audioId).UpdateByMap(map[string]interface{}{
  112. "audio_status": constvar.AudioStatusFinish,
  113. "score": resp.Score,
  114. "tags": strings.Join(words, ","),
  115. })
  116. if err != nil {
  117. return err
  118. }
  119. err = models.NewAudioTextSearch().SetOrm(db).Save(&models.AudioText{
  120. AudioID: audio.ID,
  121. AudioText: resp.Result,
  122. })
  123. return err
  124. })
  125. if err != nil {
  126. logx.Infof("AnalysisAudio success but update record failed: %v", err)
  127. _ = models.NewAudioSearch().SetID(audioId).UpdateByMap(map[string]interface{}{"audio_status": constvar.AudioStatusFailed})
  128. return
  129. }
  130. }()
  131. return nil
  132. }
  133. func GetWordFromText(text string, audio *models.Audio) (words []string) {
  134. if audio == nil {
  135. return nil
  136. }
  137. wordRecords, err := models.NewWordSearch().SetLocomotiveNumber(audio.LocomotiveNumber).FindNotTotal()
  138. if err != nil || len(wordRecords) == 0 {
  139. return nil
  140. }
  141. for _, v := range wordRecords {
  142. if strings.Contains(text, v.Content) {
  143. words = append(words, v.Content)
  144. }
  145. }
  146. return words
  147. }
  148. func PreLoad(cxt context.Context) {
  149. mkdirErr := os.MkdirAll(conf.LocalConf.PreLoadPath, os.ModePerm)
  150. if mkdirErr != nil {
  151. logx.Errorf("function os.MkdirAll() err:%v", mkdirErr)
  152. }
  153. mkdirErr1 := os.MkdirAll(conf.LocalConf.StorePath, os.ModePerm)
  154. if mkdirErr1 != nil {
  155. logx.Errorf("function os.MkdirAll() err:%v", mkdirErr1)
  156. }
  157. //文件夹下新增音频文件时触发
  158. watcher, err := fsnotify.NewWatcher()
  159. if err != nil {
  160. log.Fatal(err)
  161. }
  162. defer watcher.Close()
  163. err = watcher.Add(conf.LocalConf.PreLoadPath)
  164. if err != nil {
  165. log.Fatal(err)
  166. }
  167. audoF := func(eventName, fileName string, audio *models.Audio) bool {
  168. time.Sleep(time.Second * 1)
  169. //设置文件访问权限
  170. err = os.Chmod(eventName, 0777)
  171. if err != nil {
  172. logx.Errorf(fmt.Sprintf("%s:%s", eventName, "设置文件权限失败"))
  173. }
  174. //校验文件命名
  175. arr := strings.Split(fileName, "_")
  176. if len(arr) != 6 {
  177. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "文件名称错误"))
  178. return false
  179. }
  180. timeStr := arr[4] + strings.Split(arr[5], ".")[0]
  181. t, err := time.ParseInLocation("20060102150405", timeStr, time.Local)
  182. if err != nil {
  183. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "时间格式不对"))
  184. }
  185. //查重
  186. _, err = models.NewAudioSearch().SetName(fileName).First()
  187. if err != gorm.ErrRecordNotFound {
  188. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "重复上传"))
  189. return false
  190. }
  191. //将文件移动到uploads文件夹下
  192. //判断storePath中末尾是否带
  193. var src string
  194. if strings.HasSuffix(conf.LocalConf.StorePath, "/") {
  195. src = conf.LocalConf.StorePath + fileName
  196. } else {
  197. src = conf.LocalConf.StorePath + "/" + fileName
  198. }
  199. //err = os.Rename(eventName, src)
  200. //利用exec命令移动文件
  201. cmd := exec.Command("mv", eventName, src)
  202. err = cmd.Run()
  203. if err != nil {
  204. logx.Errorf(fmt.Sprintf("%s:%s-%s", fileName, "移动文件失败", err.Error()))
  205. return false
  206. }
  207. // 读取文件大小
  208. fileInfo, err := os.Stat(src)
  209. if err != nil {
  210. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "获取文件大小失败"))
  211. return false
  212. }
  213. size := fileInfo.Size()
  214. fmt.Println("fileName:", fileName, "size:", size, "src1", src)
  215. audio.Name = fileName
  216. audio.Size = size
  217. audio.FilePath = src
  218. audio.AudioStatus = constvar.AudioStatusUploadOk
  219. audio.LocomotiveNumber = arr[0]
  220. audio.TrainNumber = arr[1]
  221. audio.DriverNumber = arr[2]
  222. audio.Station = arr[3]
  223. audio.OccurrenceAt = t
  224. audio.IsFollowed = 0
  225. return true
  226. }
  227. txtF := func(filePath string, audio *models.Audio) bool {
  228. fileName := filepath.Base(filePath)
  229. //读取filepath文件内容到bts
  230. bts, err := os.ReadFile(filePath)
  231. if err != nil {
  232. logx.Errorf(fmt.Sprintf("%s:%s", filePath, "读取txt文件失败"))
  233. return false
  234. }
  235. //解析 交路号:123_公里标:321
  236. fileds := string(bts)
  237. arr := strings.Split(fileds, "\n")
  238. if len(arr) != 2 {
  239. logx.Errorf(fmt.Sprintf("%s:%s", filePath, "读取txt文件内容格式不对"))
  240. return false
  241. } else {
  242. RouteNumber := strings.Split(arr[0], ":")
  243. KilometerMarker := strings.Split(arr[1], ":")
  244. if len(RouteNumber) > 1 && len(KilometerMarker) > 1 {
  245. audio.RouteNumber = RouteNumber[1]
  246. audio.KilometerMarker = KilometerMarker[1]
  247. } else {
  248. logx.Errorf(fmt.Sprintf("%s:%s", filePath, "文件内容格式不对"))
  249. return false
  250. }
  251. }
  252. var src string
  253. if strings.HasSuffix(conf.LocalConf.StorePath, "/") {
  254. src = conf.LocalConf.StorePath + fileName
  255. } else {
  256. src = conf.LocalConf.StorePath + "/" + fileName
  257. }
  258. //err = os.Rename(filePath, src)
  259. //利用exec命令移动文件
  260. cmd := exec.Command("mv", filePath, src)
  261. err = cmd.Run()
  262. if err != nil {
  263. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "移动文件失败"))
  264. return false
  265. }
  266. audio.TxtFilePath = src
  267. return true
  268. }
  269. //成对变量
  270. pair := make(map[string]string)
  271. FOR:
  272. for {
  273. select {
  274. case <-cxt.Done():
  275. fmt.Println("preload stop")
  276. break FOR // 退出循环
  277. case event, ok := <-watcher.Events:
  278. if !ok {
  279. continue
  280. }
  281. if event.Op&fsnotify.Create == fsnotify.Create {
  282. // 文件名
  283. fileName := filepath.Base(event.Name)
  284. //获取不带扩展名的文件名
  285. name := strings.TrimSuffix(fileName, filepath.Ext(fileName))
  286. //判断文件在pair中
  287. if _, ok := pair[name]; !ok {
  288. pair[name] = event.Name
  289. } else {
  290. audio := &models.Audio{}
  291. isOk := true
  292. // 判断文件类型是否为.mp3或.wav
  293. if strings.ToLower(filepath.Ext(event.Name)) == ".mp3" || strings.ToLower(filepath.Ext(event.Name)) == ".wav" {
  294. isOk = audoF(event.Name, fileName, audio) && txtF(pair[name], audio)
  295. }
  296. if strings.ToLower(filepath.Ext(event.Name)) == ".txt" {
  297. isOk = audoF(pair[name], filepath.Base(pair[name]), audio) && txtF(event.Name, audio)
  298. }
  299. if !isOk {
  300. delete(pair, name)
  301. continue
  302. }
  303. if len(audio.Name) > 0 {
  304. if err = models.NewAudioSearch().Create(audio); err != nil {
  305. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "数据库create失败"))
  306. continue
  307. }
  308. go func() {
  309. var trainInfoNames = []string{audio.LocomotiveNumber, audio.TrainNumber, audio.Station} //
  310. var (
  311. info *models.TrainInfo
  312. err error
  313. parent models.TrainInfo
  314. )
  315. for i := 0; i < 3; i++ {
  316. name := trainInfoNames[i]
  317. class := constvar.Class(i + 1)
  318. info, err = models.NewTrainInfoSearch().SetName(name).SetClass(class).First()
  319. if err == gorm.ErrRecordNotFound {
  320. info = &models.TrainInfo{
  321. Name: name,
  322. Class: class,
  323. ParentID: parent.ID,
  324. }
  325. _ = models.NewTrainInfoSearch().Create(info)
  326. }
  327. parent = *info
  328. }
  329. }()
  330. }
  331. }
  332. }
  333. case err, ok := <-watcher.Errors:
  334. if !ok {
  335. logx.Errorf(err.Error())
  336. }
  337. }
  338. }
  339. }