|
|
|
@ -59,94 +59,126 @@ func (sys *NotificationSys) GetPeerRPCClient(addr xnet.Host) *PeerRPCClient { |
|
|
|
|
return sys.peerRPCClientMap[addr] |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// NotificationPeerErr returns error associated for a remote peer.
|
|
|
|
|
type NotificationPeerErr struct { |
|
|
|
|
Host xnet.Host // Remote host on which the rpc call was initiated
|
|
|
|
|
Err error // Error returned by the remote peer for an rpc call
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// DeleteBucket - calls DeleteBucket RPC call on all peers.
|
|
|
|
|
func (sys *NotificationSys) DeleteBucket(bucketName string) map[xnet.Host]error { |
|
|
|
|
errors := make(map[xnet.Host]error) |
|
|
|
|
func (sys *NotificationSys) DeleteBucket(bucketName string) []NotificationPeerErr { |
|
|
|
|
errs := make([]NotificationPeerErr, len(sys.peerRPCClientMap)) |
|
|
|
|
var wg sync.WaitGroup |
|
|
|
|
idx := 0 |
|
|
|
|
for addr, client := range sys.peerRPCClientMap { |
|
|
|
|
wg.Add(1) |
|
|
|
|
go func(addr xnet.Host, client *PeerRPCClient) { |
|
|
|
|
go func(idx int, addr xnet.Host, client *PeerRPCClient) { |
|
|
|
|
defer wg.Done() |
|
|
|
|
if err := client.DeleteBucket(bucketName); err != nil { |
|
|
|
|
errors[addr] = err |
|
|
|
|
errs[idx] = NotificationPeerErr{ |
|
|
|
|
Host: addr, |
|
|
|
|
Err: err, |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
}(addr, client) |
|
|
|
|
}(idx, addr, client) |
|
|
|
|
idx++ |
|
|
|
|
} |
|
|
|
|
wg.Wait() |
|
|
|
|
|
|
|
|
|
return errors |
|
|
|
|
return errs |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// SetBucketPolicy - calls SetBucketPolicy RPC call on all peers.
|
|
|
|
|
func (sys *NotificationSys) SetBucketPolicy(bucketName string, bucketPolicy *policy.Policy) map[xnet.Host]error { |
|
|
|
|
errors := make(map[xnet.Host]error) |
|
|
|
|
func (sys *NotificationSys) SetBucketPolicy(bucketName string, bucketPolicy *policy.Policy) []NotificationPeerErr { |
|
|
|
|
errs := make([]NotificationPeerErr, len(sys.peerRPCClientMap)) |
|
|
|
|
var wg sync.WaitGroup |
|
|
|
|
idx := 0 |
|
|
|
|
for addr, client := range sys.peerRPCClientMap { |
|
|
|
|
wg.Add(1) |
|
|
|
|
go func(addr xnet.Host, client *PeerRPCClient) { |
|
|
|
|
go func(idx int, addr xnet.Host, client *PeerRPCClient) { |
|
|
|
|
defer wg.Done() |
|
|
|
|
if err := client.SetBucketPolicy(bucketName, bucketPolicy); err != nil { |
|
|
|
|
errors[addr] = err |
|
|
|
|
errs[idx] = NotificationPeerErr{ |
|
|
|
|
Host: addr, |
|
|
|
|
Err: err, |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
}(addr, client) |
|
|
|
|
}(idx, addr, client) |
|
|
|
|
idx++ |
|
|
|
|
} |
|
|
|
|
wg.Wait() |
|
|
|
|
|
|
|
|
|
return errors |
|
|
|
|
return errs |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// RemoveBucketPolicy - calls RemoveBucketPolicy RPC call on all peers.
|
|
|
|
|
func (sys *NotificationSys) RemoveBucketPolicy(bucketName string) map[xnet.Host]error { |
|
|
|
|
errors := make(map[xnet.Host]error) |
|
|
|
|
func (sys *NotificationSys) RemoveBucketPolicy(bucketName string) []NotificationPeerErr { |
|
|
|
|
errs := make([]NotificationPeerErr, len(sys.peerRPCClientMap)) |
|
|
|
|
var wg sync.WaitGroup |
|
|
|
|
idx := 0 |
|
|
|
|
for addr, client := range sys.peerRPCClientMap { |
|
|
|
|
wg.Add(1) |
|
|
|
|
go func(addr xnet.Host, client *PeerRPCClient) { |
|
|
|
|
go func(idx int, addr xnet.Host, client *PeerRPCClient) { |
|
|
|
|
defer wg.Done() |
|
|
|
|
if err := client.RemoveBucketPolicy(bucketName); err != nil { |
|
|
|
|
errors[addr] = err |
|
|
|
|
errs[idx] = NotificationPeerErr{ |
|
|
|
|
Host: addr, |
|
|
|
|
Err: err, |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
}(addr, client) |
|
|
|
|
}(idx, addr, client) |
|
|
|
|
idx++ |
|
|
|
|
} |
|
|
|
|
wg.Wait() |
|
|
|
|
|
|
|
|
|
return errors |
|
|
|
|
return errs |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// PutBucketNotification - calls PutBucketNotification RPC call on all peers.
|
|
|
|
|
func (sys *NotificationSys) PutBucketNotification(bucketName string, rulesMap event.RulesMap) map[xnet.Host]error { |
|
|
|
|
errors := make(map[xnet.Host]error) |
|
|
|
|
func (sys *NotificationSys) PutBucketNotification(bucketName string, rulesMap event.RulesMap) []NotificationPeerErr { |
|
|
|
|
errs := make([]NotificationPeerErr, len(sys.peerRPCClientMap)) |
|
|
|
|
var wg sync.WaitGroup |
|
|
|
|
idx := 0 |
|
|
|
|
for addr, client := range sys.peerRPCClientMap { |
|
|
|
|
wg.Add(1) |
|
|
|
|
go func(addr xnet.Host, client *PeerRPCClient, rulesMap event.RulesMap) { |
|
|
|
|
go func(idx int, addr xnet.Host, client *PeerRPCClient, rulesMap event.RulesMap) { |
|
|
|
|
defer wg.Done() |
|
|
|
|
if err := client.PutBucketNotification(bucketName, rulesMap); err != nil { |
|
|
|
|
errors[addr] = err |
|
|
|
|
errs[idx] = NotificationPeerErr{ |
|
|
|
|
Host: addr, |
|
|
|
|
Err: err, |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
}(addr, client, rulesMap.Clone()) |
|
|
|
|
}(idx, addr, client, rulesMap.Clone()) |
|
|
|
|
idx++ |
|
|
|
|
} |
|
|
|
|
wg.Wait() |
|
|
|
|
|
|
|
|
|
return errors |
|
|
|
|
return errs |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// ListenBucketNotification - calls ListenBucketNotification RPC call on all peers.
|
|
|
|
|
func (sys *NotificationSys) ListenBucketNotification(bucketName string, eventNames []event.Name, pattern string, targetID event.TargetID, localPeer xnet.Host) map[xnet.Host]error { |
|
|
|
|
errors := make(map[xnet.Host]error) |
|
|
|
|
func (sys *NotificationSys) ListenBucketNotification(bucketName string, eventNames []event.Name, pattern string, |
|
|
|
|
targetID event.TargetID, localPeer xnet.Host) []NotificationPeerErr { |
|
|
|
|
errs := make([]NotificationPeerErr, len(sys.peerRPCClientMap)) |
|
|
|
|
var wg sync.WaitGroup |
|
|
|
|
idx := 0 |
|
|
|
|
for addr, client := range sys.peerRPCClientMap { |
|
|
|
|
wg.Add(1) |
|
|
|
|
go func(addr xnet.Host, client *PeerRPCClient) { |
|
|
|
|
go func(idx int, addr xnet.Host, client *PeerRPCClient) { |
|
|
|
|
defer wg.Done() |
|
|
|
|
if err := client.ListenBucketNotification(bucketName, eventNames, pattern, targetID, localPeer); err != nil { |
|
|
|
|
errors[addr] = err |
|
|
|
|
errs[idx] = NotificationPeerErr{ |
|
|
|
|
Host: addr, |
|
|
|
|
Err: err, |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
}(addr, client) |
|
|
|
|
}(idx, addr, client) |
|
|
|
|
idx++ |
|
|
|
|
} |
|
|
|
|
wg.Wait() |
|
|
|
|
|
|
|
|
|
return errors |
|
|
|
|
return errs |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// AddRemoteTarget - adds event rules map, HTTP/PeerRPC client target to bucket name.
|
|
|
|
|