aboutsummaryrefslogtreecommitdiff
path: root/weed/filesys/page_writer/upload_pipeline_lock.go
blob: 47a40ba3771e89e8728d97879c239aeedf07c122 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
package page_writer

import (
	"sync/atomic"
)

func (up *UploadPipeline) LockForRead(startOffset, stopOffset int64) {
	startLogicChunkIndex := LogicChunkIndex(startOffset / up.ChunkSize)
	stopLogicChunkIndex := LogicChunkIndex(stopOffset / up.ChunkSize)
	if stopOffset%up.ChunkSize > 0 {
		stopLogicChunkIndex += 1
	}
	up.activeReadChunksLock.Lock()
	defer up.activeReadChunksLock.Unlock()
	for i := startLogicChunkIndex; i < stopLogicChunkIndex; i++ {
		if count, found := up.activeReadChunks[i]; found {
			up.activeReadChunks[i] = count + 1
		} else {
			up.activeReadChunks[i] = 1
		}
	}
}

func (up *UploadPipeline) UnlockForRead(startOffset, stopOffset int64) {
	startLogicChunkIndex := LogicChunkIndex(startOffset / up.ChunkSize)
	stopLogicChunkIndex := LogicChunkIndex(stopOffset / up.ChunkSize)
	if stopOffset%up.ChunkSize > 0 {
		stopLogicChunkIndex += 1
	}
	up.activeReadChunksLock.Lock()
	defer up.activeReadChunksLock.Unlock()
	for i := startLogicChunkIndex; i < stopLogicChunkIndex; i++ {
		if count, found := up.activeReadChunks[i]; found {
			if count == 1 {
				delete(up.activeReadChunks, i)
			} else {
				up.activeReadChunks[i] = count - 1
			}
		}
	}
}

func (up *UploadPipeline) IsLocked(logicChunkIndex LogicChunkIndex) bool {
	up.activeReadChunksLock.Lock()
	defer up.activeReadChunksLock.Unlock()
	if count, found := up.activeReadChunks[logicChunkIndex]; found {
		return count > 0
	}
	return false
}

func (up *UploadPipeline) waitForCurrentWritersToComplete() {
	up.uploaderCountCond.L.Lock()
	t := int32(100)
	for {
		t = atomic.LoadInt32(&up.uploaderCount)
		if t <= 0 {
			break
		}
		up.uploaderCountCond.Wait()
	}
	up.uploaderCountCond.L.Unlock()
}