package service import ( "encoding/json" "sync" "time" "3588AdminBackend/internal/config" "3588AdminBackend/internal/models" ) type RegistryService struct { cfg *config.Config agent *AgentClient mu sync.RWMutex devices map[string]*models.Device } func NewRegistryService(cfg *config.Config, agent *AgentClient) *RegistryService { s := &RegistryService{ cfg: cfg, agent: agent, devices: make(map[string]*models.Device), } go s.startPruning() go s.startGraphPolling() return s } func (s *RegistryService) startGraphPolling() { ticker := time.NewTicker(30 * time.Second) // Pull every 30s for range ticker.C { s.mu.RLock() var onlineDevices []*models.Device for _, dev := range s.devices { if dev.Online { onlineDevices = append(onlineDevices, dev) } } s.mu.RUnlock() for _, dev := range onlineDevices { data, _, err := s.agent.Do("GET", dev.IP, dev.AgentPort, "/v1/graphs", nil) if err == nil { var graphs interface{} if err := json.Unmarshal(data, &graphs); err == nil { s.mu.Lock() dev.Graphs = graphs s.mu.Unlock() } } } } } func (s *RegistryService) UpdateDevice(dev *models.Device) { s.mu.Lock() defer s.mu.Unlock() dev.LastSeenMs = time.Now().UnixMilli() dev.Online = true s.devices[dev.DeviceID] = dev } func (s *RegistryService) GetDevices() []*models.Device { s.mu.RLock() defer s.mu.RUnlock() list := make([]*models.Device, 0, len(s.devices)) for _, dev := range s.devices { list = append(list, dev) } return list } func (s *RegistryService) startPruning() { ticker := time.NewTicker(2 * time.Second) for range ticker.C { s.mu.Lock() now := time.Now().UnixMilli() for _, dev := range s.devices { if now-dev.LastSeenMs > int64(s.cfg.OfflineAfterMs) { dev.Online = false } } s.mu.Unlock() } }