forked from vitessio/vitess
-
Notifications
You must be signed in to change notification settings - Fork 0
/
resource_constraint.go
106 lines (87 loc) · 2.53 KB
/
resource_constraint.go
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
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
// Copyright 2013, Google Inc. All rights reserved.
// Use of this source code is governed by a BSD-style
// license that can be found in the LICENSE file.
package concurrency
import (
"fmt"
"sync"
"github.com/youtube/vitess/go/sync2"
)
// ResourceConstraint combines 3 different features:
// - a WaitGroup to wait for all tasks to be done
// - a Semaphore to control concurrency
// - an ErrorRecorder
type ResourceConstraint struct {
semaphore *sync2.Semaphore
wg sync.WaitGroup
FirstErrorRecorder
}
// NewResourceConstraint creates a ResourceConstraint with
// max concurrency.
func NewResourceConstraint(max int) *ResourceConstraint {
return &ResourceConstraint{semaphore: sync2.NewSemaphore(max, 0)}
}
func (rc *ResourceConstraint) Add(n int) {
rc.wg.Add(n)
}
func (rc *ResourceConstraint) Done() {
rc.wg.Done()
}
// Wait waits for the WG and returns the firstError we encountered, or nil
func (rc *ResourceConstraint) Wait() error {
rc.wg.Wait()
return rc.Error()
}
// Acquire will wait until we have a resource to use
func (rc *ResourceConstraint) Acquire() {
rc.semaphore.Acquire()
}
func (rc *ResourceConstraint) Release() {
rc.semaphore.Release()
}
func (rc *ResourceConstraint) ReleaseAndDone() {
rc.Release()
rc.Done()
}
// MultiResourceConstraint combines 3 different features:
// - a WaitGroup to wait for all tasks to be done
// - a Semaphore map to control multiple concurrencies
// - an ErrorRecorder
type MultiResourceConstraint struct {
semaphoreMap map[string]*sync2.Semaphore
wg sync.WaitGroup
FirstErrorRecorder
}
func NewMultiResourceConstraint(semaphoreMap map[string]*sync2.Semaphore) *MultiResourceConstraint {
return &MultiResourceConstraint{semaphoreMap: semaphoreMap}
}
func (mrc *MultiResourceConstraint) Add(n int) {
mrc.wg.Add(n)
}
func (mrc *MultiResourceConstraint) Done() {
mrc.wg.Done()
}
// Returns the firstError we encountered, or nil
func (mrc *MultiResourceConstraint) Wait() error {
mrc.wg.Wait()
return mrc.Error()
}
// Acquire will wait until we have a resource to use
func (mrc *MultiResourceConstraint) Acquire(name string) {
s, ok := mrc.semaphoreMap[name]
if !ok {
panic(fmt.Errorf("MultiResourceConstraint: No resource named %v in semaphore map", name))
}
s.Acquire()
}
func (mrc *MultiResourceConstraint) Release(name string) {
s, ok := mrc.semaphoreMap[name]
if !ok {
panic(fmt.Errorf("MultiResourceConstraint: No resource named %v in semaphore map", name))
}
s.Release()
}
func (mrc *MultiResourceConstraint) ReleaseAndDone(name string) {
mrc.Release(name)
mrc.Done()
}