...
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15 package v2v3
16
17 import (
18 "context"
19 "strings"
20
21 "go.etcd.io/etcd/client/v3"
22 "go.etcd.io/etcd/server/v3/etcdserver/api/v2error"
23 "go.etcd.io/etcd/server/v3/etcdserver/api/v2store"
24 )
25
26 func (s *v2v3Store) Watch(prefix string, recursive, stream bool, sinceIndex uint64) (v2store.Watcher, error) {
27 ctx, cancel := context.WithCancel(s.ctx)
28 wch := s.c.Watch(
29 ctx,
30
31 s.pfx,
32 clientv3.WithPrefix(),
33 clientv3.WithRev(int64(sinceIndex)),
34 clientv3.WithCreatedNotify(),
35 clientv3.WithPrevKV())
36 resp, ok := <-wch
37 if err := resp.Err(); err != nil || !ok {
38 cancel()
39 return nil, v2error.NewError(v2error.EcodeRaftInternal, prefix, 0)
40 }
41
42 evc, donec := make(chan *v2store.Event), make(chan struct{})
43 go func() {
44 defer func() {
45 close(evc)
46 close(donec)
47 }()
48 for resp := range wch {
49 for _, ev := range s.mkV2Events(resp) {
50 k := ev.Node.Key
51 if recursive {
52 if !strings.HasPrefix(k, prefix) {
53 continue
54 }
55
56 k = strings.Replace(k, prefix, "/", 1)
57
58 if strings.Contains(k, "/_") {
59 continue
60 }
61 }
62 if !recursive && k != prefix {
63 continue
64 }
65 select {
66 case evc <- ev:
67 case <-ctx.Done():
68 return
69 }
70 if !stream {
71 return
72 }
73 }
74 }
75 }()
76
77 return &v2v3Watcher{
78 startRev: resp.Header.Revision,
79 evc: evc,
80 donec: donec,
81 cancel: cancel,
82 }, nil
83 }
84
85 func (s *v2v3Store) mkV2Events(wr clientv3.WatchResponse) (evs []*v2store.Event) {
86 ak := s.mkActionKey()
87 for _, rev := range mkRevs(wr) {
88 var act, key *clientv3.Event
89 for _, ev := range rev {
90 if string(ev.Kv.Key) == ak {
91 act = ev
92 } else if key != nil && len(key.Kv.Key) < len(ev.Kv.Key) {
93
94
95 key = ev
96 } else if key == nil {
97 key = ev
98 }
99 }
100 if act != nil && act.Kv != nil && key != nil {
101 v2ev := &v2store.Event{
102 Action: string(act.Kv.Value),
103 Node: s.mkV2Node(key.Kv),
104 PrevNode: s.mkV2Node(key.PrevKv),
105 EtcdIndex: mkV2Rev(wr.Header.Revision),
106 }
107 evs = append(evs, v2ev)
108 }
109 }
110 return evs
111 }
112
113 func mkRevs(wr clientv3.WatchResponse) (revs [][]*clientv3.Event) {
114 var curRev []*clientv3.Event
115 for _, ev := range wr.Events {
116 if curRev != nil && ev.Kv.ModRevision != curRev[0].Kv.ModRevision {
117 revs = append(revs, curRev)
118 curRev = nil
119 }
120 curRev = append(curRev, ev)
121 }
122 if curRev != nil {
123 revs = append(revs, curRev)
124 }
125 return revs
126 }
127
128 type v2v3Watcher struct {
129 startRev int64
130 evc chan *v2store.Event
131 donec chan struct{}
132 cancel context.CancelFunc
133 }
134
135 func (w *v2v3Watcher) StartIndex() uint64 { return mkV2Rev(w.startRev) }
136
137 func (w *v2v3Watcher) Remove() {
138 w.cancel()
139 <-w.donec
140 }
141
142 func (w *v2v3Watcher) EventChan() chan *v2store.Event { return w.evc }
143
View as plain text