-
Notifications
You must be signed in to change notification settings - Fork 1
/
notify.go
69 lines (53 loc) · 1.18 KB
/
notify.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
package klevdb
import (
"context"
"sync/atomic"
"github.com/klev-dev/kleverr"
)
type OffsetNotify struct {
nextOffset atomic.Int64
barrier chan chan struct{}
}
func NewOffsetNotify(nextOffset int64) *OffsetNotify {
w := &OffsetNotify{
barrier: make(chan chan struct{}, 1),
}
w.nextOffset.Store(nextOffset)
w.barrier <- make(chan struct{})
return w
}
func (w *OffsetNotify) Wait(ctx context.Context, offset int64) error {
// quick path, just load and check
if w.nextOffset.Load() > offset {
return nil
}
// acquire current barrier
b := <-w.barrier
// probe the current offset
updated := w.nextOffset.Load() > offset
// release current barrier
w.barrier <- b
// already has a new value, return
if updated {
return nil
}
// now wait for something to happen
select {
case <-b:
return nil
case <-ctx.Done():
return kleverr.Ret(ctx.Err())
}
}
func (w *OffsetNotify) Set(nextOffset int64) {
// acquire current barrier
b := <-w.barrier
// set the new offset
if w.nextOffset.Load() < nextOffset {
w.nextOffset.Store(nextOffset)
}
// close the current barrier, e.g. broadcast update
close(b)
// create new barrier
w.barrier <- make(chan struct{})
}