collector_go/collector.go

134 lines
2.6 KiB
Go
Raw Normal View History

2017-04-13 03:43:39 +00:00
package collector_go
import (
2017-05-11 09:34:47 +00:00
"context"
2017-05-15 10:32:48 +00:00
cm "loafle.com/overflow/agent_api/config_manager"
2017-05-11 09:34:47 +00:00
"loafle.com/overflow/crawler_go/grpc"
crm "loafle.com/overflow/crawler_manager_go"
s "loafle.com/overflow/scheduler_go"
2017-04-13 03:43:39 +00:00
"log"
2017-05-11 09:34:47 +00:00
"strconv"
2017-04-28 04:16:08 +00:00
"sync"
2017-04-15 11:32:56 +00:00
"time"
2017-04-13 03:43:39 +00:00
)
2017-04-28 04:16:08 +00:00
var (
instance *Collector
once sync.Once
)
2017-04-13 03:43:39 +00:00
2017-05-15 10:32:48 +00:00
type Collector struct {
scheduler *s.Scheduler
cm cm.ConfigManager
2017-05-16 07:44:02 +00:00
dataCh chan interface{}
2017-05-15 10:32:48 +00:00
}
func Start(started chan bool, dataCh chan interface{}, conf cm.ConfigManager) {
c := GetInstance()
c.dataCh = dataCh
c.start(started, conf)
}
func Stop(stopped chan bool) {
c := GetInstance()
c.stop()
stopped <- true
2017-04-14 10:09:51 +00:00
}
2017-04-13 03:43:39 +00:00
2017-04-28 04:16:08 +00:00
func GetInstance() *Collector {
once.Do(func() {
instance = &Collector{}
})
return instance
}
2017-04-15 11:32:56 +00:00
2017-05-15 10:32:48 +00:00
func (c *Collector) start(started chan bool, conf cm.ConfigManager) {
2017-04-28 04:16:08 +00:00
go func() {
c.cm = conf
2017-05-11 09:34:47 +00:00
c.scheduler = &s.Scheduler{}
c.scheduler.Start()
2017-04-15 11:32:56 +00:00
2017-04-28 04:16:08 +00:00
for _, conf := range c.cm.GetSensors() {
if err := c.addSensor(conf.Id); err != nil {
2017-04-15 11:32:56 +00:00
log.Println(err)
}
}
2017-05-15 10:32:48 +00:00
started <- true
2017-04-15 11:32:56 +00:00
}()
}
2017-05-15 10:32:48 +00:00
func (c *Collector) stop() {
2017-04-14 10:09:51 +00:00
c.scheduler.RemoveAllSchedule()
2017-04-15 11:32:56 +00:00
c.scheduler.Stop()
2017-04-14 10:09:51 +00:00
}
2017-04-13 03:43:39 +00:00
2017-04-28 04:16:08 +00:00
func (c *Collector) collect(id string) {
2017-04-13 03:43:39 +00:00
2017-04-28 04:16:08 +00:00
conf := c.cm.GetSensorById(id)
log.Printf("COLLECT %s - [ID: %s] [Crawler : %s]", time.Now(), conf.Id, conf.Crawler.Name)
2017-05-11 09:34:47 +00:00
conn, err := crm.GetInstance().GetClient(conf.Crawler.Container)
if err != nil {
log.Println(err)
}
defer conn.Close()
2017-04-15 11:32:56 +00:00
2017-05-11 09:34:47 +00:00
dc := grpc.NewDataClient(conn)
in := &grpc.Input{}
2017-04-15 11:32:56 +00:00
2017-05-11 09:34:47 +00:00
in.Id = id
in.Name = grpc.Crawlers(grpc.Crawlers_value[conf.Crawler.Name])
2017-04-15 11:32:56 +00:00
2017-05-11 09:34:47 +00:00
out, err := dc.Get(context.Background(), in)
2017-04-15 11:32:56 +00:00
2017-05-11 09:34:47 +00:00
if err != nil {
log.Println(err)
}
2017-05-15 10:32:48 +00:00
log.Println("collector get result : ", out)
c.dataCh <- out
2017-04-15 11:32:56 +00:00
}
2017-04-28 04:16:08 +00:00
func (c *Collector) addSensor(sensorId string) error {
sensor := c.cm.GetSensorById(sensorId)
2017-05-11 09:34:47 +00:00
interval, err := strconv.Atoi(sensor.Schedule.Interval)
if err != nil {
return err
}
return c.scheduler.NewSchedule(sensorId, uint64(interval), c.collect)
2017-04-14 10:09:51 +00:00
}
2017-04-13 03:43:39 +00:00
2017-05-15 10:32:48 +00:00
func (c *Collector) removeSensor(id string) error {
2017-04-28 04:16:08 +00:00
if err := c.scheduler.RemoveSchedule(id); err != nil {
2017-05-15 10:32:48 +00:00
return err
2017-04-13 03:43:39 +00:00
}
2017-05-15 10:32:48 +00:00
return nil
2017-04-14 10:09:51 +00:00
}
2017-05-11 09:34:47 +00:00
2017-05-16 07:44:02 +00:00
func (c *Collector) updateSensor(id string) error {
err := c.removeSensor(id)
if err != nil {
return err
}
return c.addSensor(id)
2017-05-11 09:34:47 +00:00
}
2017-05-15 10:32:48 +00:00
func AddSensor(id string) error {
return GetInstance().addSensor(id)
}
func RemSensor(id string) error {
return GetInstance().removeSensor(id)
}
2017-05-16 07:44:02 +00:00
func UpdateSensor(id string) error {
return GetInstance().updateSensor(id)
}
2017-05-15 10:32:48 +00:00
func StartSensor(id string) error {
return GetInstance().scheduler.StartSchedule(id)
}
func StopSensor(id string) error {
return GetInstance().scheduler.StopSchedule(id)
2017-05-16 07:44:02 +00:00
}