代码拉取完成,页面将自动刷新
package harvester
import (
"sync"
uuid "github.com/satori/go.uuid"
"github.com/elastic/beats/libbeat/logp"
)
// Registry struct manages (start / stop) a list of harvesters
type Registry struct {
sync.RWMutex
harvesters map[uuid.UUID]Harvester
wg sync.WaitGroup
done chan struct{}
}
// NewRegistry creates a new registry object
func NewRegistry() *Registry {
return &Registry{
harvesters: map[uuid.UUID]Harvester{},
done: make(chan struct{}),
}
}
func (r *Registry) remove(h Harvester) {
r.Lock()
defer r.Unlock()
delete(r.harvesters, h.ID())
}
func (r *Registry) add(h Harvester) {
r.Lock()
defer r.Unlock()
r.harvesters[h.ID()] = h
}
// Stop stops all harvesters in the registry
func (r *Registry) Stop() {
r.Lock()
defer func() {
r.Unlock()
r.WaitForCompletion()
}()
// Makes sure no new harvesters are added during stopping
close(r.done)
for _, hv := range r.harvesters {
go func(h Harvester) {
h.Stop()
}(hv)
}
}
// WaitForCompletion can be used to wait until all harvesters are stopped
func (r *Registry) WaitForCompletion() {
r.wg.Wait()
}
// Start starts the given harvester and add its to the registry
func (r *Registry) Start(h Harvester) {
// Make sure stop is not called during starting a harvester
r.Lock()
defer r.Unlock()
// Make sure no new harvesters are started after stop was called
if !r.active() {
return
}
r.wg.Add(1)
go func() {
defer func() {
r.remove(h)
r.wg.Done()
}()
r.add(h)
// Starts harvester and picks the right type. In case type is not set, set it to default (log)
err := h.Run()
if err != nil {
logp.Err("Error running prospector: %v", err)
}
}()
}
// Len returns the current number of harvesters in the registry
func (r *Registry) Len() uint64 {
r.RLock()
defer r.RUnlock()
return uint64(len(r.harvesters))
}
func (r *Registry) active() bool {
select {
case <-r.done:
return false
default:
return true
}
}
此处可能存在不合适展示的内容,页面不予展示。您可通过相关编辑功能自查并修改。
如您确认内容无涉及 不当用语 / 纯广告导流 / 暴力 / 低俗色情 / 侵权 / 盗版 / 虚假 / 无价值内容或违法国家有关法律法规的内容,可点击提交进行申诉,我们将尽快为您处理。