You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

queue_disk_channel_test.go 11KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544
  1. // Copyright 2019 The Gitea Authors. All rights reserved.
  2. // SPDX-License-Identifier: MIT
  3. package queue
  4. import (
  5. "sync"
  6. "testing"
  7. "time"
  8. "code.gitea.io/gitea/modules/log"
  9. "github.com/stretchr/testify/assert"
  10. )
  11. func TestPersistableChannelQueue(t *testing.T) {
  12. handleChan := make(chan *testData)
  13. handle := func(data ...Data) []Data {
  14. for _, datum := range data {
  15. if datum == nil {
  16. continue
  17. }
  18. testDatum := datum.(*testData)
  19. handleChan <- testDatum
  20. }
  21. return nil
  22. }
  23. lock := sync.Mutex{}
  24. queueShutdown := []func(){}
  25. queueTerminate := []func(){}
  26. tmpDir := t.TempDir()
  27. queue, err := NewPersistableChannelQueue(handle, PersistableChannelQueueConfiguration{
  28. DataDir: tmpDir,
  29. BatchLength: 2,
  30. QueueLength: 20,
  31. Workers: 1,
  32. BoostWorkers: 0,
  33. MaxWorkers: 10,
  34. Name: "first",
  35. }, &testData{})
  36. assert.NoError(t, err)
  37. readyForShutdown := make(chan struct{})
  38. readyForTerminate := make(chan struct{})
  39. go queue.Run(func(shutdown func()) {
  40. lock.Lock()
  41. defer lock.Unlock()
  42. select {
  43. case <-readyForShutdown:
  44. default:
  45. close(readyForShutdown)
  46. }
  47. queueShutdown = append(queueShutdown, shutdown)
  48. }, func(terminate func()) {
  49. lock.Lock()
  50. defer lock.Unlock()
  51. select {
  52. case <-readyForTerminate:
  53. default:
  54. close(readyForTerminate)
  55. }
  56. queueTerminate = append(queueTerminate, terminate)
  57. })
  58. test1 := testData{"A", 1}
  59. test2 := testData{"B", 2}
  60. err = queue.Push(&test1)
  61. assert.NoError(t, err)
  62. go func() {
  63. err := queue.Push(&test2)
  64. assert.NoError(t, err)
  65. }()
  66. result1 := <-handleChan
  67. assert.Equal(t, test1.TestString, result1.TestString)
  68. assert.Equal(t, test1.TestInt, result1.TestInt)
  69. result2 := <-handleChan
  70. assert.Equal(t, test2.TestString, result2.TestString)
  71. assert.Equal(t, test2.TestInt, result2.TestInt)
  72. // test1 is a testData not a *testData so will be rejected
  73. err = queue.Push(test1)
  74. assert.Error(t, err)
  75. <-readyForShutdown
  76. // Now shutdown the queue
  77. lock.Lock()
  78. callbacks := make([]func(), len(queueShutdown))
  79. copy(callbacks, queueShutdown)
  80. lock.Unlock()
  81. for _, callback := range callbacks {
  82. callback()
  83. }
  84. // Wait til it is closed
  85. <-queue.(*PersistableChannelQueue).closed
  86. err = queue.Push(&test1)
  87. assert.NoError(t, err)
  88. err = queue.Push(&test2)
  89. assert.NoError(t, err)
  90. select {
  91. case <-handleChan:
  92. assert.Fail(t, "Handler processing should have stopped")
  93. default:
  94. }
  95. // terminate the queue
  96. <-readyForTerminate
  97. lock.Lock()
  98. callbacks = make([]func(), len(queueTerminate))
  99. copy(callbacks, queueTerminate)
  100. lock.Unlock()
  101. for _, callback := range callbacks {
  102. callback()
  103. }
  104. select {
  105. case <-handleChan:
  106. assert.Fail(t, "Handler processing should have stopped")
  107. default:
  108. }
  109. // Reopen queue
  110. queue, err = NewPersistableChannelQueue(handle, PersistableChannelQueueConfiguration{
  111. DataDir: tmpDir,
  112. BatchLength: 2,
  113. QueueLength: 20,
  114. Workers: 1,
  115. BoostWorkers: 0,
  116. MaxWorkers: 10,
  117. Name: "second",
  118. }, &testData{})
  119. assert.NoError(t, err)
  120. readyForShutdown = make(chan struct{})
  121. readyForTerminate = make(chan struct{})
  122. go queue.Run(func(shutdown func()) {
  123. lock.Lock()
  124. defer lock.Unlock()
  125. select {
  126. case <-readyForShutdown:
  127. default:
  128. close(readyForShutdown)
  129. }
  130. queueShutdown = append(queueShutdown, shutdown)
  131. }, func(terminate func()) {
  132. lock.Lock()
  133. defer lock.Unlock()
  134. select {
  135. case <-readyForTerminate:
  136. default:
  137. close(readyForTerminate)
  138. }
  139. queueTerminate = append(queueTerminate, terminate)
  140. })
  141. result3 := <-handleChan
  142. assert.Equal(t, test1.TestString, result3.TestString)
  143. assert.Equal(t, test1.TestInt, result3.TestInt)
  144. result4 := <-handleChan
  145. assert.Equal(t, test2.TestString, result4.TestString)
  146. assert.Equal(t, test2.TestInt, result4.TestInt)
  147. <-readyForShutdown
  148. lock.Lock()
  149. callbacks = make([]func(), len(queueShutdown))
  150. copy(callbacks, queueShutdown)
  151. lock.Unlock()
  152. for _, callback := range callbacks {
  153. callback()
  154. }
  155. <-readyForTerminate
  156. lock.Lock()
  157. callbacks = make([]func(), len(queueTerminate))
  158. copy(callbacks, queueTerminate)
  159. lock.Unlock()
  160. for _, callback := range callbacks {
  161. callback()
  162. }
  163. }
  164. func TestPersistableChannelQueue_Pause(t *testing.T) {
  165. lock := sync.Mutex{}
  166. var queue Queue
  167. var err error
  168. pushBack := false
  169. handleChan := make(chan *testData)
  170. handle := func(data ...Data) []Data {
  171. lock.Lock()
  172. if pushBack {
  173. if pausable, ok := queue.(Pausable); ok {
  174. log.Info("pausing")
  175. pausable.Pause()
  176. }
  177. lock.Unlock()
  178. return data
  179. }
  180. lock.Unlock()
  181. for _, datum := range data {
  182. testDatum := datum.(*testData)
  183. handleChan <- testDatum
  184. }
  185. return nil
  186. }
  187. queueShutdown := []func(){}
  188. queueTerminate := []func(){}
  189. terminated := make(chan struct{})
  190. tmpDir := t.TempDir()
  191. queue, err = NewPersistableChannelQueue(handle, PersistableChannelQueueConfiguration{
  192. DataDir: tmpDir,
  193. BatchLength: 2,
  194. QueueLength: 20,
  195. Workers: 1,
  196. BoostWorkers: 0,
  197. MaxWorkers: 10,
  198. Name: "first",
  199. }, &testData{})
  200. assert.NoError(t, err)
  201. go func() {
  202. queue.Run(func(shutdown func()) {
  203. lock.Lock()
  204. defer lock.Unlock()
  205. queueShutdown = append(queueShutdown, shutdown)
  206. }, func(terminate func()) {
  207. lock.Lock()
  208. defer lock.Unlock()
  209. queueTerminate = append(queueTerminate, terminate)
  210. })
  211. close(terminated)
  212. }()
  213. // Shutdown and Terminate in defer
  214. defer func() {
  215. lock.Lock()
  216. callbacks := make([]func(), len(queueShutdown))
  217. copy(callbacks, queueShutdown)
  218. lock.Unlock()
  219. for _, callback := range callbacks {
  220. callback()
  221. }
  222. lock.Lock()
  223. log.Info("Finally terminating")
  224. callbacks = make([]func(), len(queueTerminate))
  225. copy(callbacks, queueTerminate)
  226. lock.Unlock()
  227. for _, callback := range callbacks {
  228. callback()
  229. }
  230. }()
  231. test1 := testData{"A", 1}
  232. test2 := testData{"B", 2}
  233. err = queue.Push(&test1)
  234. assert.NoError(t, err)
  235. pausable, ok := queue.(Pausable)
  236. if !assert.True(t, ok) {
  237. return
  238. }
  239. result1 := <-handleChan
  240. assert.Equal(t, test1.TestString, result1.TestString)
  241. assert.Equal(t, test1.TestInt, result1.TestInt)
  242. pausable.Pause()
  243. paused, _ := pausable.IsPausedIsResumed()
  244. select {
  245. case <-paused:
  246. case <-time.After(100 * time.Millisecond):
  247. assert.Fail(t, "Queue is not paused")
  248. return
  249. }
  250. queue.Push(&test2)
  251. var result2 *testData
  252. select {
  253. case result2 = <-handleChan:
  254. assert.Fail(t, "handler chan should be empty")
  255. case <-time.After(100 * time.Millisecond):
  256. }
  257. assert.Nil(t, result2)
  258. pausable.Resume()
  259. _, resumed := pausable.IsPausedIsResumed()
  260. select {
  261. case <-resumed:
  262. case <-time.After(100 * time.Millisecond):
  263. assert.Fail(t, "Queue should be resumed")
  264. return
  265. }
  266. select {
  267. case result2 = <-handleChan:
  268. case <-time.After(500 * time.Millisecond):
  269. assert.Fail(t, "handler chan should contain test2")
  270. }
  271. assert.Equal(t, test2.TestString, result2.TestString)
  272. assert.Equal(t, test2.TestInt, result2.TestInt)
  273. // Set pushBack to so that the next handle will result in a Pause
  274. lock.Lock()
  275. pushBack = true
  276. lock.Unlock()
  277. // Ensure that we're still resumed
  278. _, resumed = pausable.IsPausedIsResumed()
  279. select {
  280. case <-resumed:
  281. case <-time.After(100 * time.Millisecond):
  282. assert.Fail(t, "Queue is not resumed")
  283. return
  284. }
  285. // push test1
  286. queue.Push(&test1)
  287. // Now as this is handled it should pause
  288. paused, _ = pausable.IsPausedIsResumed()
  289. select {
  290. case <-paused:
  291. case <-handleChan:
  292. assert.Fail(t, "handler chan should not contain test1")
  293. return
  294. case <-time.After(500 * time.Millisecond):
  295. assert.Fail(t, "queue should be paused")
  296. return
  297. }
  298. lock.Lock()
  299. pushBack = false
  300. lock.Unlock()
  301. pausable.Resume()
  302. _, resumed = pausable.IsPausedIsResumed()
  303. select {
  304. case <-resumed:
  305. case <-time.After(500 * time.Millisecond):
  306. assert.Fail(t, "Queue should be resumed")
  307. return
  308. }
  309. select {
  310. case result1 = <-handleChan:
  311. case <-time.After(500 * time.Millisecond):
  312. assert.Fail(t, "handler chan should contain test1")
  313. return
  314. }
  315. assert.Equal(t, test1.TestString, result1.TestString)
  316. assert.Equal(t, test1.TestInt, result1.TestInt)
  317. lock.Lock()
  318. callbacks := make([]func(), len(queueShutdown))
  319. copy(callbacks, queueShutdown)
  320. queueShutdown = queueShutdown[:0]
  321. lock.Unlock()
  322. // Now shutdown the queue
  323. for _, callback := range callbacks {
  324. callback()
  325. }
  326. // Wait til it is closed
  327. select {
  328. case <-queue.(*PersistableChannelQueue).closed:
  329. case <-time.After(5 * time.Second):
  330. assert.Fail(t, "queue should close")
  331. return
  332. }
  333. err = queue.Push(&test1)
  334. assert.NoError(t, err)
  335. err = queue.Push(&test2)
  336. assert.NoError(t, err)
  337. select {
  338. case <-handleChan:
  339. assert.Fail(t, "Handler processing should have stopped")
  340. return
  341. default:
  342. }
  343. // terminate the queue
  344. lock.Lock()
  345. callbacks = make([]func(), len(queueTerminate))
  346. copy(callbacks, queueTerminate)
  347. queueShutdown = queueTerminate[:0]
  348. lock.Unlock()
  349. for _, callback := range callbacks {
  350. callback()
  351. }
  352. select {
  353. case <-handleChan:
  354. assert.Fail(t, "Handler processing should have stopped")
  355. return
  356. case <-terminated:
  357. case <-time.After(10 * time.Second):
  358. assert.Fail(t, "Queue should have terminated")
  359. return
  360. }
  361. lock.Lock()
  362. pushBack = true
  363. lock.Unlock()
  364. // Reopen queue
  365. terminated = make(chan struct{})
  366. queue, err = NewPersistableChannelQueue(handle, PersistableChannelQueueConfiguration{
  367. DataDir: tmpDir,
  368. BatchLength: 1,
  369. QueueLength: 20,
  370. Workers: 1,
  371. BoostWorkers: 0,
  372. MaxWorkers: 10,
  373. Name: "second",
  374. }, &testData{})
  375. assert.NoError(t, err)
  376. pausable, ok = queue.(Pausable)
  377. if !assert.True(t, ok) {
  378. return
  379. }
  380. paused, _ = pausable.IsPausedIsResumed()
  381. go func() {
  382. queue.Run(func(shutdown func()) {
  383. lock.Lock()
  384. defer lock.Unlock()
  385. queueShutdown = append(queueShutdown, shutdown)
  386. }, func(terminate func()) {
  387. lock.Lock()
  388. defer lock.Unlock()
  389. queueTerminate = append(queueTerminate, terminate)
  390. })
  391. close(terminated)
  392. }()
  393. select {
  394. case <-handleChan:
  395. assert.Fail(t, "Handler processing should have stopped")
  396. return
  397. case <-paused:
  398. }
  399. paused, _ = pausable.IsPausedIsResumed()
  400. select {
  401. case <-paused:
  402. case <-time.After(500 * time.Millisecond):
  403. assert.Fail(t, "Queue is not paused")
  404. return
  405. }
  406. select {
  407. case <-handleChan:
  408. assert.Fail(t, "Handler processing should have stopped")
  409. return
  410. default:
  411. }
  412. lock.Lock()
  413. pushBack = false
  414. lock.Unlock()
  415. pausable.Resume()
  416. _, resumed = pausable.IsPausedIsResumed()
  417. select {
  418. case <-resumed:
  419. case <-time.After(500 * time.Millisecond):
  420. assert.Fail(t, "Queue should be resumed")
  421. return
  422. }
  423. var result3, result4 *testData
  424. select {
  425. case result3 = <-handleChan:
  426. case <-time.After(1 * time.Second):
  427. assert.Fail(t, "Handler processing should have resumed")
  428. return
  429. }
  430. select {
  431. case result4 = <-handleChan:
  432. case <-time.After(1 * time.Second):
  433. assert.Fail(t, "Handler processing should have resumed")
  434. return
  435. }
  436. if result4.TestString == test1.TestString {
  437. result3, result4 = result4, result3
  438. }
  439. assert.Equal(t, test1.TestString, result3.TestString)
  440. assert.Equal(t, test1.TestInt, result3.TestInt)
  441. assert.Equal(t, test2.TestString, result4.TestString)
  442. assert.Equal(t, test2.TestInt, result4.TestInt)
  443. lock.Lock()
  444. callbacks = make([]func(), len(queueShutdown))
  445. copy(callbacks, queueShutdown)
  446. queueShutdown = queueShutdown[:0]
  447. lock.Unlock()
  448. // Now shutdown the queue
  449. for _, callback := range callbacks {
  450. callback()
  451. }
  452. // terminate the queue
  453. lock.Lock()
  454. callbacks = make([]func(), len(queueTerminate))
  455. copy(callbacks, queueTerminate)
  456. queueShutdown = queueTerminate[:0]
  457. lock.Unlock()
  458. for _, callback := range callbacks {
  459. callback()
  460. }
  461. select {
  462. case <-time.After(10 * time.Second):
  463. assert.Fail(t, "Queue should have terminated")
  464. return
  465. case <-terminated:
  466. }
  467. }