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.

57 lines
1.4 KiB

3 years ago
  1. package page_writer
  2. import (
  3. "sync/atomic"
  4. )
  5. func (up *UploadPipeline) LockForRead(startOffset, stopOffset int64) {
  6. startLogicChunkIndex := LogicChunkIndex(startOffset / up.ChunkSize)
  7. stopLogicChunkIndex := LogicChunkIndex(stopOffset / up.ChunkSize)
  8. if stopOffset%up.ChunkSize > 0 {
  9. stopLogicChunkIndex += 1
  10. }
  11. for i := startLogicChunkIndex; i < stopLogicChunkIndex; i++ {
  12. if count, found := up.activeReadChunks[i]; found {
  13. up.activeReadChunks[i] = count + 1
  14. } else {
  15. up.activeReadChunks[i] = 1
  16. }
  17. }
  18. }
  19. func (up *UploadPipeline) UnlockForRead(startOffset, stopOffset int64) {
  20. startLogicChunkIndex := LogicChunkIndex(startOffset / up.ChunkSize)
  21. stopLogicChunkIndex := LogicChunkIndex(stopOffset / up.ChunkSize)
  22. if stopOffset%up.ChunkSize > 0 {
  23. stopLogicChunkIndex += 1
  24. }
  25. for i := startLogicChunkIndex; i < stopLogicChunkIndex; i++ {
  26. if count, found := up.activeReadChunks[i]; found {
  27. if count == 1 {
  28. delete(up.activeReadChunks, i)
  29. } else {
  30. up.activeReadChunks[i] = count - 1
  31. }
  32. }
  33. }
  34. }
  35. func (up *UploadPipeline) IsLocked(logicChunkIndex LogicChunkIndex) bool {
  36. if count, found := up.activeReadChunks[logicChunkIndex]; found {
  37. return count > 0
  38. }
  39. return false
  40. }
  41. func (up *UploadPipeline) waitForCurrentWritersToComplete() {
  42. up.uploaderCountCond.L.Lock()
  43. t := int32(100)
  44. for {
  45. t = atomic.LoadInt32(&up.uploaderCount)
  46. if t <= 0 {
  47. break
  48. }
  49. up.uploaderCountCond.Wait()
  50. }
  51. up.uploaderCountCond.L.Unlock()
  52. }