process.go 9.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362
  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. var resp Response
  91. sz := audio.Size / 1024 / 1024
  92. if sz > 2 {
  93. resp, err = AnalysisAudio(audio.FilePath, conf.AanlysisConf.LongUrl)
  94. } else {
  95. resp, err = AnalysisAudio(audio.FilePath, conf.AanlysisConf.Url)
  96. }
  97. if err != nil {
  98. logx.Errorf("err when AnalysisAudio:%v", err)
  99. _ = models.NewAudioSearch().SetID(audioId).UpdateByMap(map[string]interface{}{"audio_status": constvar.AudioStatusFailed})
  100. return
  101. }
  102. if resp.Code != 0 {
  103. logx.Errorf("AnalysisAudio error return:%v", resp)
  104. _ = models.NewAudioSearch().SetID(audioId).UpdateByMap(map[string]interface{}{"audio_status": constvar.AudioStatusFailed})
  105. return
  106. }
  107. logx.Infof("AnalysisAudio result: %v", resp)
  108. words := GetWordFromText(resp.Result, audio)
  109. err = models.WithTransaction(func(db *gorm.DB) error {
  110. err = models.NewAudioSearch().SetOrm(db).SetID(audioId).UpdateByMap(map[string]interface{}{
  111. "audio_status": constvar.AudioStatusFinish,
  112. "score": resp.Score,
  113. "tags": strings.Join(words, ","),
  114. })
  115. if err != nil {
  116. return err
  117. }
  118. err = models.NewAudioTextSearch().SetOrm(db).Save(&models.AudioText{
  119. AudioID: audio.ID,
  120. AudioText: resp.Result,
  121. })
  122. return err
  123. })
  124. if err != nil {
  125. logx.Infof("AnalysisAudio success but update record failed: %v", err)
  126. _ = models.NewAudioSearch().SetID(audioId).UpdateByMap(map[string]interface{}{"audio_status": constvar.AudioStatusFailed})
  127. return
  128. }
  129. }()
  130. return nil
  131. }
  132. func GetWordFromText(text string, audio *models.Audio) (words []string) {
  133. if audio == nil {
  134. return nil
  135. }
  136. wordRecords, err := models.NewWordSearch().SetLocomotiveNumber(audio.LocomotiveNumber).FindNotTotal()
  137. if err != nil || len(wordRecords) == 0 {
  138. return nil
  139. }
  140. for _, v := range wordRecords {
  141. if strings.Contains(text, v.Content) {
  142. words = append(words, v.Content)
  143. }
  144. }
  145. return words
  146. }
  147. func PreLoad(cxt context.Context) {
  148. mkdirErr := os.MkdirAll(conf.LocalConf.PreLoadPath, os.ModePerm)
  149. if mkdirErr != nil {
  150. logx.Errorf("function os.MkdirAll() err:%v", mkdirErr)
  151. }
  152. mkdirErr1 := os.MkdirAll(conf.LocalConf.StorePath, os.ModePerm)
  153. if mkdirErr1 != nil {
  154. logx.Errorf("function os.MkdirAll() err:%v", mkdirErr1)
  155. }
  156. //文件夹下新增音频文件时触发
  157. watcher, err := fsnotify.NewWatcher()
  158. if err != nil {
  159. log.Fatal(err)
  160. }
  161. defer watcher.Close()
  162. err = watcher.Add(conf.LocalConf.PreLoadPath)
  163. if err != nil {
  164. log.Fatal(err)
  165. }
  166. audoF := func(eventName, fileName string, audio *models.Audio) bool {
  167. time.Sleep(time.Second * 1)
  168. //设置文件访问权限
  169. err = os.Chmod(eventName, 0777)
  170. if err != nil {
  171. logx.Errorf(fmt.Sprintf("%s:%s", eventName, "设置文件权限失败"))
  172. }
  173. //校验文件命名
  174. arr := strings.Split(fileName, "_")
  175. if len(arr) != 6 {
  176. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "文件名称错误"))
  177. return false
  178. }
  179. timeStr := arr[4] + strings.Split(arr[5], ".")[0]
  180. t, err := time.ParseInLocation("20060102150405", timeStr, time.Local)
  181. if err != nil {
  182. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "时间格式不对"))
  183. }
  184. //查重
  185. _, err = models.NewAudioSearch().SetName(fileName).First()
  186. if err != gorm.ErrRecordNotFound {
  187. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "重复上传"))
  188. return false
  189. }
  190. //将文件移动到uploads文件夹下
  191. //判断storePath中末尾是否带
  192. var src string
  193. if strings.HasSuffix(conf.LocalConf.StorePath, "/") {
  194. src = conf.LocalConf.StorePath + fileName
  195. } else {
  196. src = conf.LocalConf.StorePath + "/" + fileName
  197. }
  198. err = os.Rename(eventName, src)
  199. if err != nil {
  200. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "移动文件失败"))
  201. return false
  202. }
  203. // 读取文件大小
  204. fileInfo, err := os.Stat(src)
  205. if err != nil {
  206. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "获取文件大小失败"))
  207. return false
  208. }
  209. size := fileInfo.Size()
  210. fmt.Println("fileName:", fileName, "size:", size, "src1", src)
  211. audio.Name = fileName
  212. audio.Size = size
  213. audio.FilePath = src
  214. audio.AudioStatus = constvar.AudioStatusUploadOk
  215. audio.LocomotiveNumber = arr[0]
  216. audio.TrainNumber = arr[1]
  217. audio.DriverNumber = arr[2]
  218. audio.Station = arr[3]
  219. audio.OccurrenceAt = t
  220. audio.IsFollowed = 0
  221. return true
  222. }
  223. txtF := func(filePath string, audio *models.Audio) bool {
  224. fileName := filepath.Base(filePath)
  225. //读取filepath文件内容到bts
  226. bts, err := os.ReadFile(filePath)
  227. if err != nil {
  228. logx.Errorf(fmt.Sprintf("%s:%s", filePath, "读取txt文件失败"))
  229. return false
  230. }
  231. //解析 交路号:123_公里标:321
  232. fileds := string(bts)
  233. arr := strings.Split(fileds, "_")
  234. if len(arr) != 2 {
  235. logx.Errorf(fmt.Sprintf("%s:%s", filePath, "读取txt文件内容格式不对"))
  236. return false
  237. } else {
  238. RouteNumber := strings.Split(arr[0], ":")
  239. KilometerMarker := strings.Split(arr[1], ":")
  240. if len(RouteNumber) > 1 && len(KilometerMarker) > 1 {
  241. audio.RouteNumber = RouteNumber[1]
  242. audio.KilometerMarker = KilometerMarker[1]
  243. } else {
  244. logx.Errorf(fmt.Sprintf("%s:%s", filePath, "文件内容格式不对"))
  245. return false
  246. }
  247. }
  248. var src string
  249. if strings.HasSuffix(conf.LocalConf.StorePath, "/") {
  250. src = conf.LocalConf.StorePath + fileName
  251. } else {
  252. src = conf.LocalConf.StorePath + "/" + fileName
  253. }
  254. err = os.Rename(filePath, src)
  255. if err != nil {
  256. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "移动文件失败"))
  257. return false
  258. }
  259. audio.TxtFilePath = src
  260. return true
  261. }
  262. //成对变量
  263. pair := make(map[string]string)
  264. FOR:
  265. for {
  266. select {
  267. case <-cxt.Done():
  268. fmt.Println("preload stop")
  269. break FOR // 退出循环
  270. case event, ok := <-watcher.Events:
  271. if !ok {
  272. continue
  273. }
  274. if event.Op&fsnotify.Create == fsnotify.Create {
  275. // 文件名
  276. fileName := filepath.Base(event.Name)
  277. //获取不带扩展名的文件名
  278. name := strings.TrimSuffix(fileName, filepath.Ext(fileName))
  279. //判断文件在pair中
  280. if _, ok := pair[name]; !ok {
  281. pair[name] = event.Name
  282. } else {
  283. audio := &models.Audio{}
  284. isOk := true
  285. // 判断文件类型是否为.mp3或.wav
  286. if filepath.Ext(event.Name) == ".mp3" || filepath.Ext(event.Name) == ".wav" {
  287. isOk = audoF(event.Name, fileName, audio) && txtF(pair[name], audio)
  288. }
  289. if filepath.Ext(event.Name) == ".txt" {
  290. isOk = audoF(pair[name], filepath.Base(pair[name]), audio) && txtF(event.Name, audio)
  291. }
  292. if !isOk {
  293. delete(pair, name)
  294. continue
  295. }
  296. if len(audio.Name) > 0 {
  297. if err = models.NewAudioSearch().Create(audio); err != nil {
  298. logx.Errorf(fmt.Sprintf("%s:%s", fileName, "数据库create失败"))
  299. continue
  300. }
  301. go func() {
  302. var trainInfoNames = []string{audio.LocomotiveNumber, audio.TrainNumber, audio.Station} //
  303. var (
  304. info *models.TrainInfo
  305. err error
  306. parent models.TrainInfo
  307. )
  308. for i := 0; i < 3; i++ {
  309. name := trainInfoNames[i]
  310. class := constvar.Class(i + 1)
  311. info, err = models.NewTrainInfoSearch().SetName(name).SetClass(class).First()
  312. if err == gorm.ErrRecordNotFound {
  313. info = &models.TrainInfo{
  314. Name: name,
  315. Class: class,
  316. ParentID: parent.ID,
  317. }
  318. _ = models.NewTrainInfoSearch().Create(info)
  319. }
  320. parent = *info
  321. }
  322. }()
  323. }
  324. }
  325. }
  326. case err, ok := <-watcher.Errors:
  327. if !ok {
  328. logx.Errorf(err.Error())
  329. }
  330. }
  331. }
  332. }