mql_run_file.go 7.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288
  1. package odbcmql
  2. import (
  3. "bufio"
  4. "context"
  5. "fmt"
  6. "os"
  7. "path/filepath"
  8. "strings"
  9. "sync"
  10. "sync/atomic"
  11. "testing"
  12. "time"
  13. "gitee.com/wecisecode/util/pqc"
  14. "gitee.com/wecisecode/util/spliter"
  15. "github.com/stretchr/testify/assert"
  16. )
  17. func (mt *MQLTest) RunFile(t *testing.T, ctx context.Context,
  18. global *GlobalVars,
  19. topvars *CurrentVars,
  20. dirvars *CurrentVars,
  21. basedir, filename string) bool {
  22. // 读取文件内容
  23. ffpath := filepath.Join(basedir, filename)
  24. f, err := os.Open(ffpath)
  25. if !assert.Nil(t, err, err) {
  26. return false
  27. }
  28. // mql语句切分
  29. mqgr := NewMQLGroupRequest()
  30. var multilines *MQLRequest
  31. mqs := spliter.NewMQLSpliter(bufio.NewReader(f))
  32. rtinfo := time.Now().Add(5 * time.Second)
  33. for {
  34. mql, fromline, toline, fromchar, tochar, hasnext, _ := mqs.NextMQL()
  35. if !hasnext {
  36. break
  37. }
  38. if time.Now().After(rtinfo) {
  39. logger.Info("reading bigfile ", toline, " lines")
  40. rtinfo = time.Now().Add(5 * time.Second)
  41. }
  42. // 去掉注释
  43. clean_mql := strings.Join(spliter.MQLSplitClean(mql), ";")
  44. if multilines != nil {
  45. // 多条语句处理
  46. if strings.TrimSpace(clean_mql) == "multilines end" {
  47. // 保留语句中的注释
  48. tm := strings.TrimSpace(strings.Replace(mql, "multilines end", "", 1))
  49. if len(tm) > 0 {
  50. if multilines.OriginQueryString != "" {
  51. multilines.OriginQueryString += ";"
  52. }
  53. multilines.OriginQueryString += tm
  54. }
  55. multilines.Toline = toline
  56. multilines.Tochar = tochar
  57. // 加入请求组,并初始化,提取动作信息
  58. e := mqgr.Append(multilines)
  59. if !assert.Nil(t, e, e) {
  60. return false
  61. }
  62. multilines = nil
  63. } else {
  64. if multilines.OriginQueryString != "" {
  65. multilines.OriginQueryString += ";"
  66. }
  67. multilines.OriginQueryString += mql
  68. multilines.Toline = toline
  69. multilines.Tochar = tochar
  70. }
  71. } else if strings.TrimSpace(clean_mql) == "multilines begin" {
  72. multilines = &MQLRequest{FilePath: ffpath}
  73. // 保留语句中的注释
  74. multilines.OriginQueryString = strings.TrimSpace(strings.Replace(mql, "multilines begin", "", 1))
  75. multilines.Fromline = fromline
  76. multilines.Fromchar = fromchar
  77. multilines.Toline = toline
  78. multilines.Tochar = tochar
  79. } else {
  80. // 单条语句
  81. mqr := &MQLRequest{OriginQueryString: mql, FilePath: ffpath, Fromline: fromline, Toline: toline, Fromchar: fromchar, Tochar: tochar}
  82. // 加入请求组,并初始化,提取动作信息
  83. e := mqgr.Append(mqr)
  84. if !assert.Nil(t, e, e) {
  85. return false
  86. }
  87. }
  88. }
  89. // mqls := spliter.MQLSplit(string(bs))
  90. mt.scopevars.Lock()
  91. if mt.scopevars.file[ffpath] == nil {
  92. mt.scopevars.file[ffpath] = &Variables{
  93. vars: map[string]interface{}{},
  94. loop_count: 1,
  95. loop_from: 1,
  96. loop_step: 1}
  97. }
  98. mt.scopevars.Unlock()
  99. var wg sync.WaitGroup
  100. st := time.Now()
  101. loop_i := 0
  102. parallel_queue := pqc.NewQueue[any](0)
  103. mqlcount := int32(0)
  104. for {
  105. mt.scopevars.Lock()
  106. ok := loop_i < mt.scopevars.file[ffpath].loop_count
  107. mt.scopevars.Unlock()
  108. if !ok {
  109. break
  110. }
  111. loop_i++
  112. //
  113. ch_parallel_count := make(chan int)
  114. filevars := &CurrentVars{
  115. loop_i: loop_i,
  116. ch_parallel_count: ch_parallel_count,
  117. }
  118. //
  119. ok_chan := make(chan bool, 1)
  120. parallel_chan := make(chan bool, 1)
  121. parallel := false
  122. parallelcount := 0
  123. done := false
  124. wg.Add(1)
  125. go func() {
  126. defer func() {
  127. atomic.AddInt32(&mqlcount, filevars.mqlcount)
  128. wg.Done()
  129. }()
  130. ch_ok := make(chan bool)
  131. go func() {
  132. for {
  133. select {
  134. case <-ch_parallel_count:
  135. if !done && !parallel {
  136. parallel = true
  137. // 加入并发控制队列
  138. if parallelcount > 0 {
  139. if parallelcount > parallel_queue.Size() {
  140. parallel_queue.Growth(parallelcount)
  141. }
  142. parallel_queue.Push(1)
  143. }
  144. parallel_chan <- true
  145. }
  146. case ok := <-ch_ok:
  147. ok_chan <- ok
  148. if parallel {
  149. if parallelcount > 0 {
  150. // 从并发控制队列中移除
  151. parallel_queue.Pop()
  152. }
  153. } else {
  154. done = true
  155. }
  156. return
  157. }
  158. }
  159. }()
  160. // logger.Info("file", ffpath, "第", fmt.Sprint(topvars.loop_i, ".", dirvars.loop_i, ".", filevars.loop_i), "次执行开始")
  161. ch_ok <- mt.RunMQLGroup(t, ctx,
  162. global,
  163. topvars,
  164. dirvars,
  165. filevars,
  166. basedir, ffpath, mqgr)
  167. // logger.Info("file", ffpath, "第", fmt.Sprint(topvars.loop_i, ".", dirvars.loop_i, ".", filevars.loop_i), "次执行结束")
  168. }()
  169. success := true
  170. select {
  171. case success = <-ok_chan: // 非并发,等待完成
  172. logger.Info("file", ffpath, "第", fmt.Sprint(topvars.loop_i, ".", dirvars.loop_i, ".", filevars.loop_i), "次顺序执行完成")
  173. case <-parallel_chan: // 并发,执行继续下一次
  174. logger.Info("file", ffpath, "第", fmt.Sprint(topvars.loop_i, ".", dirvars.loop_i, ".", filevars.loop_i), "次并发执行继续")
  175. }
  176. if !success {
  177. return false
  178. }
  179. }
  180. wg.Wait()
  181. mt.scopevars.RLock()
  182. loop_count := mt.scopevars.file[ffpath].loop_count
  183. mt.scopevars.RUnlock()
  184. if loop_count > 1 {
  185. ut := time.Since(st)
  186. sn := fmt.Sprint(topvars.loop_i, ".", dirvars.loop_i)
  187. logger.Info(fmt.Sprint("file ", ffpath+"/"+sn, " loop ", loop_count, " times, run ", mqlcount, " mqls, usetime ", ut))
  188. }
  189. return true
  190. }
  191. func (mt *MQLTest) RunMQLGroup(t *testing.T, ctx context.Context,
  192. global *GlobalVars,
  193. topvars *CurrentVars,
  194. dirvars *CurrentVars,
  195. filevars *CurrentVars,
  196. basedir, ffpath string, mqgr *MQLGroupRequest) (pass bool) {
  197. pass = true
  198. var wg sync.WaitGroup
  199. for _, mqs := range mqgr.mqrs {
  200. wg.Add(1)
  201. go func(mqs []*MQLRequest) {
  202. defer wg.Done()
  203. if len(mqs) == 0 {
  204. return
  205. }
  206. if mqs[0].StaticActions.ForkName != nil {
  207. forkname := *mqs[0].StaticActions.ForkName
  208. global.Lock()
  209. wg := global.wg_wait_fork_routine[forkname]
  210. if wg == nil {
  211. wg = &sync.WaitGroup{}
  212. global.wg_wait_fork_routine[forkname] = wg
  213. }
  214. global.Unlock()
  215. wg.Add(1)
  216. defer wg.Done()
  217. }
  218. ok := mt.RunMQLs(t, ctx, global, topvars, dirvars, filevars, basedir, ffpath, mqs)
  219. if !ok {
  220. pass = false
  221. }
  222. }(mqs)
  223. }
  224. wg.Wait()
  225. return
  226. }
  227. func (mt *MQLTest) RunMQLs(t *testing.T, ctx context.Context,
  228. global *GlobalVars,
  229. topvars *CurrentVars,
  230. dirvars *CurrentVars,
  231. filevars *CurrentVars,
  232. basedir, ffpath string, mqrs []*MQLRequest) bool {
  233. for _, mqr := range mqrs {
  234. mqlkey := mqr.Key
  235. mqlstr := mqr.OriginQueryString
  236. if mqlstr == "" {
  237. continue
  238. }
  239. staticactions := mqr.StaticActions
  240. // 设置执行过程中的控制参数
  241. mt.InitScopeVars(basedir, ffpath, mqlkey, staticactions)
  242. ch_test_run_one_mql_result := make(chan bool)
  243. go func() {
  244. mqrinst := fmt.Sprint(mqlkey, "/", topvars.loop_i, ".", dirvars.loop_i, ".", filevars.loop_i)
  245. global.Lock()
  246. mqrdone := global.ch_wait_mql_done[mqrinst]
  247. if mqrdone == nil {
  248. mqrdone = make(chan bool)
  249. global.ch_wait_mql_done[mqrinst] = mqrdone
  250. }
  251. var waitmqrdone chan bool
  252. if mqr.WaitMQLRequest != nil {
  253. waitmqrkey := mqr.WaitMQLRequest.Key
  254. waitmqrinst := fmt.Sprint(waitmqrkey, "/", topvars.loop_i, ".", dirvars.loop_i, ".", filevars.loop_i)
  255. waitmqrdone = global.ch_wait_mql_done[waitmqrinst]
  256. if waitmqrdone == nil {
  257. waitmqrdone = make(chan bool)
  258. global.ch_wait_mql_done[waitmqrinst] = waitmqrdone
  259. }
  260. }
  261. global.Unlock()
  262. var ret bool
  263. defer func() { mqrdone <- ret }()
  264. if waitmqrdone != nil {
  265. v := <-waitmqrdone
  266. waitmqrdone <- v
  267. if !v {
  268. // 依赖MQR失败
  269. ch_test_run_one_mql_result <- false
  270. return
  271. }
  272. }
  273. ret = mt.RunMQR(t, ctx, global, topvars, dirvars, filevars, basedir, ffpath, mqr)
  274. ch_test_run_one_mql_result <- ret
  275. }()
  276. ret := <-ch_test_run_one_mql_result
  277. if !ret {
  278. return ret
  279. }
  280. // continue
  281. }
  282. return true
  283. }