diff --git a/Makefile b/Makefile index efbc20f..91aaf90 100644 --- a/Makefile +++ b/Makefile @@ -1,5 +1,5 @@ BINARY=syncserver -VERSION?=1.0.10 +VERSION?=1.0.11 GO?=go LDFLAGS=-s -w -X main.version=$(VERSION) -X main.commit=$(shell git rev-parse --short HEAD 2>/dev/null || echo unknown) BUILD_FLAGS=CGO_ENABLED=0 diff --git a/cmd/server/main.go b/cmd/server/main.go index e692697..91ed5d4 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -20,7 +20,7 @@ import ( "github.com/syncserver/internal/syncengine" ) -var version = "1.0.10" +var version = "1.0.11" func main() { cfgPath := flag.String("config", "", "Path to config.yaml") diff --git a/internal/api/handlers_machines.go b/internal/api/handlers_machines.go index 8d398a3..ca14648 100644 --- a/internal/api/handlers_machines.go +++ b/internal/api/handlers_machines.go @@ -10,15 +10,17 @@ import ( "github.com/go-chi/chi/v5" "github.com/syncserver/internal/models" + "github.com/syncserver/internal/syncengine" "github.com/syncserver/internal/wol" ) type MachineHandler struct { - db *sql.DB + db *sql.DB + engine *syncengine.Engine } -func NewMachineHandler(db *sql.DB) *MachineHandler { - return &MachineHandler{db: db} +func NewMachineHandler(db *sql.DB, engine *syncengine.Engine) *MachineHandler { + return &MachineHandler{db: db, engine: engine} } var macRegex = regexp.MustCompile(`^([0-9A-Fa-f]{2}[:-]){5}[0-9A-Fa-f]{2}$`) @@ -209,6 +211,16 @@ func (h *MachineHandler) TestWoL(w http.ResponseWriter, r *http.Request) { writeJSON(w, map[string]interface{}{"ok": true, "sent": 3}) } +func (h *MachineHandler) Refresh(w http.ResponseWriter, r *http.Request) { + if h.engine == nil { + writeError(w, http.StatusInternalServerError, "engine not available") + return + } + go h.engine.ProbeAllMachines() + w.WriteHeader(http.StatusAccepted) + writeJSON(w, map[string]string{"status": "probing"}) +} + func (h *MachineHandler) Delete(w http.ResponseWriter, r *http.Request) { id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64) if err != nil { diff --git a/internal/api/router.go b/internal/api/router.go index cbd3c1e..7b691f1 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -33,7 +33,7 @@ func NewServer(cfg *config.Config, db *sql.DB, engine *syncengine.Engine) *Serve s := &Server{router: r, cfg: cfg, engine: engine} authHandler := NewAuthHandler(db) - machineHandler := NewMachineHandler(db) + machineHandler := NewMachineHandler(db, engine) syncPairHandler := NewSyncPairHandler(db) jobHandler := NewJobHandler(db, engine) sseHandler := NewSSEHandler(engine) @@ -49,6 +49,7 @@ func NewServer(cfg *config.Config, db *sql.DB, engine *syncengine.Engine) *Serve r.With(auth.RequireAuth).Route("/machines", func(r chi.Router) { r.Get("/", machineHandler.List) r.Post("/", machineHandler.Create) + r.Post("/refresh", machineHandler.Refresh) r.Get("/{id}", machineHandler.Get) r.Put("/{id}", machineHandler.Update) r.Delete("/{id}", machineHandler.Delete) diff --git a/internal/syncengine/engine.go b/internal/syncengine/engine.go index 45c3a2a..b4a4915 100644 --- a/internal/syncengine/engine.go +++ b/internal/syncengine/engine.go @@ -17,12 +17,13 @@ import ( ) type Engine struct { - db *sql.DB - cfg *config.Config - queue *Queue - eventBus *EventBus - mu sync.RWMutex - stopped bool + db *sql.DB + cfg *config.Config + queue *Queue + eventBus *EventBus + mu sync.RWMutex + stopped bool + lastProbeAt atomic.Int64 } type Event struct { @@ -333,3 +334,51 @@ func (e *Engine) CreateJob(syncPairID int64, triggerType string) (int64, error) } return id, nil } + +const ( + probeThrottleSeconds = 10 + probeTimeout = 1500 * time.Millisecond + probeMaxConcurrent = 20 +) + +func (e *Engine) ProbeAllMachines() { + now := time.Now().UnixNano() + last := e.lastProbeAt.Load() + + if now-last < int64(probeThrottleSeconds*time.Second) { + slog.Debug("ProbeAllMachines: skipped (throttled)") + return + } + if !e.lastProbeAt.CompareAndSwap(last, now) { + return + } + + machineRepo := models.NewMachineRepository(e.db) + ms, err := machineRepo.GetAll() + if err != nil { + slog.Warn("ProbeAllMachines: list failed", "error", err) + return + } + + sem := make(chan struct{}, probeMaxConcurrent) + var wg sync.WaitGroup + + for i := range ms { + wg.Add(1) + sem <- struct{}{} + go func(m *models.Machine) { + defer wg.Done() + defer func() { <-sem }() + + ctx, cancel := context.WithTimeout(context.Background(), probeTimeout) + defer cancel() + + status := "offline" + if wol.IsReachable(ctx, m.Host, m.Port, probeTimeout) { + status = "online" + } + e.setMachineStatus(m.ID, status) + }(&ms[i]) + } + wg.Wait() +} diff --git a/internal/syncengine/engine_status_test.go b/internal/syncengine/engine_status_test.go index a4f8037..cdeca8f 100644 --- a/internal/syncengine/engine_status_test.go +++ b/internal/syncengine/engine_status_test.go @@ -3,6 +3,7 @@ package syncengine import ( "context" "net" + "sync/atomic" "testing" "time" @@ -37,3 +38,29 @@ func TestIsReachableTimeout(t *testing.T) { t.Errorf("IsReachable returned too early: %v", d) } } + +func TestProbeThrottle(t *testing.T) { + e := &Engine{lastProbeAt: atomic.Int64{}} + + e.lastProbeAt.Store(time.Now().UnixNano()) + + now := time.Now().UnixNano() + last := e.lastProbeAt.Load() + if now-last < int64(probeThrottleSeconds*time.Second) { + return + } + t.Error("throttle check did not run as expected") +} + +func TestProbeConstants(t *testing.T) { + if probeThrottleSeconds != 10 { + t.Errorf("probeThrottleSeconds = %d, want 10", probeThrottleSeconds) + } + if probeTimeout != 1500*time.Millisecond { + t.Errorf("probeTimeout = %v, want 1500ms", probeTimeout) + } + if probeMaxConcurrent != 20 { + t.Errorf("probeMaxConcurrent = %d, want 20", probeMaxConcurrent) + } +} + diff --git a/web/src/pages/Dashboard.tsx b/web/src/pages/Dashboard.tsx index 50efecd..aac268c 100644 --- a/web/src/pages/Dashboard.tsx +++ b/web/src/pages/Dashboard.tsx @@ -32,6 +32,7 @@ export default function Dashboard() { }) .catch(() => {}) .finally(() => setLoading(false)); + api('/api/machines/refresh', { method: 'POST' }).catch(() => {}); }, []); useEffect(() => { diff --git a/web/src/pages/Machines.tsx b/web/src/pages/Machines.tsx index d534e53..e8c185c 100644 --- a/web/src/pages/Machines.tsx +++ b/web/src/pages/Machines.tsx @@ -66,9 +66,14 @@ export default function Machines() { const [deleteId, setDeleteId] = useState(null); const [form, setForm] = useState(defaultForm); const [loading, setLoading] = useState(false); + const [probing, setProbing] = useState(false); useEffect(() => { load(); + setProbing(true); + api('/api/machines/refresh', { method: 'POST' }).catch(() => {}); + const timer = setTimeout(() => setProbing(false), 5000); + return () => clearTimeout(timer); }, []); useEffect(() => { @@ -212,6 +217,13 @@ export default function Machines() { } /> ) : ( + <> + {probing && ( +
+ + Checking machine status... +
+ )} @@ -290,6 +302,7 @@ export default function Machines() { ))}
+ )}