process.go 9.3 KB

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