mirror of
https://github.com/etcd-io/etcd.git
synced 2024-09-27 06:25:44 +00:00
api/v3rpc: add watch ID to "watchStream.Watch"
Signed-off-by: Gyuho Lee <gyuhox@gmail.com>
This commit is contained in:
parent
82a164e3b9
commit
33c732b97c
@ -205,7 +205,7 @@ func (sws *serverWatchStream) recvLoop() error {
|
||||
if !sws.isWatchPermitted(creq) {
|
||||
wr := &pb.WatchResponse{
|
||||
Header: sws.newResponseHeader(sws.watchStream.Rev()),
|
||||
WatchId: -1,
|
||||
WatchId: creq.WatchId,
|
||||
Canceled: true,
|
||||
Created: true,
|
||||
CancelReason: rpctypes.ErrGRPCPermissionDenied.Error(),
|
||||
@ -225,8 +225,8 @@ func (sws *serverWatchStream) recvLoop() error {
|
||||
if rev == 0 {
|
||||
rev = wsrev + 1
|
||||
}
|
||||
id := sws.watchStream.Watch(creq.Key, creq.RangeEnd, rev, filters...)
|
||||
if id != -1 {
|
||||
id, err := sws.watchStream.Watch(mvcc.WatchID(creq.WatchId), creq.Key, creq.RangeEnd, rev, filters...)
|
||||
if err == nil {
|
||||
sws.mu.Lock()
|
||||
if creq.ProgressNotify {
|
||||
sws.progress[id] = true
|
||||
@ -240,7 +240,10 @@ func (sws *serverWatchStream) recvLoop() error {
|
||||
Header: sws.newResponseHeader(wsrev),
|
||||
WatchId: int64(id),
|
||||
Created: true,
|
||||
Canceled: id == -1,
|
||||
Canceled: err != nil,
|
||||
}
|
||||
if err != nil {
|
||||
wr.CancelReason = err.Error()
|
||||
}
|
||||
select {
|
||||
case sws.ctrlStream <- wr:
|
||||
|
Loading…
x
Reference in New Issue
Block a user