287 lines
6.8 KiB
Go
287 lines
6.8 KiB
Go
package health
|
|
|
|
import (
|
|
"fmt"
|
|
"sync"
|
|
"time"
|
|
|
|
"codeberg.org/SquidSE/ClusterCanvas/service/internal/auth"
|
|
"codeberg.org/SquidSE/ClusterCanvas/service/internal/settings"
|
|
)
|
|
|
|
const schedulerTickInterval = 15 * time.Second
|
|
|
|
// Result is the outcome of a node health check.
|
|
type Result struct {
|
|
PingOK bool
|
|
SSHOK bool
|
|
OK bool
|
|
Message string
|
|
CheckedAt time.Time
|
|
}
|
|
|
|
// Checker runs and schedules per-node health checks.
|
|
type Checker struct {
|
|
ConfigDir string
|
|
Key []byte
|
|
|
|
mu sync.Mutex
|
|
running map[string]struct{}
|
|
stopCh chan struct{}
|
|
stopped chan struct{}
|
|
tickerOn bool
|
|
|
|
// Optional overrides for tests.
|
|
PingFn func(hostIP string) auth.PingResult
|
|
SSHFn func(hostIP string, username string, privateKeyPEM string, passphrase string) error
|
|
NowFn func() time.Time
|
|
}
|
|
|
|
// NewChecker creates a health checker for the given config directory and key.
|
|
func NewChecker(configDir string, key []byte) *Checker {
|
|
return &Checker{
|
|
ConfigDir: configDir,
|
|
Key: key,
|
|
running: map[string]struct{}{},
|
|
stopCh: make(chan struct{}),
|
|
stopped: make(chan struct{}),
|
|
}
|
|
}
|
|
|
|
// StartScheduler begins the background health-check ticker.
|
|
func (checker *Checker) StartScheduler() {
|
|
checker.mu.Lock()
|
|
if checker.tickerOn {
|
|
checker.mu.Unlock()
|
|
return
|
|
}
|
|
checker.tickerOn = true
|
|
checker.mu.Unlock()
|
|
|
|
go func() {
|
|
defer close(checker.stopped)
|
|
ticker := time.NewTicker(schedulerTickInterval)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-checker.stopCh:
|
|
return
|
|
case <-ticker.C:
|
|
checker.tickDueChecks()
|
|
}
|
|
}
|
|
}()
|
|
}
|
|
|
|
// StopScheduler stops the background ticker and waits for it to exit.
|
|
func (checker *Checker) StopScheduler() {
|
|
checker.mu.Lock()
|
|
if !checker.tickerOn {
|
|
checker.mu.Unlock()
|
|
return
|
|
}
|
|
checker.tickerOn = false
|
|
checker.mu.Unlock()
|
|
close(checker.stopCh)
|
|
<-checker.stopped
|
|
}
|
|
|
|
func (checker *Checker) tryLockNode(nodeID string) bool {
|
|
checker.mu.Lock()
|
|
defer checker.mu.Unlock()
|
|
if _, ok := checker.running[nodeID]; ok {
|
|
return false
|
|
}
|
|
checker.running[nodeID] = struct{}{}
|
|
return true
|
|
}
|
|
|
|
func (checker *Checker) unlockNode(nodeID string) {
|
|
checker.mu.Lock()
|
|
defer checker.mu.Unlock()
|
|
delete(checker.running, nodeID)
|
|
}
|
|
|
|
func (checker *Checker) now() time.Time {
|
|
if checker.NowFn != nil {
|
|
return checker.NowFn()
|
|
}
|
|
return time.Now().UTC()
|
|
}
|
|
|
|
func (checker *Checker) tickDueChecks() {
|
|
if len(checker.Key) == 0 {
|
|
return
|
|
}
|
|
store, err := settings.LoadNodesOrEmpty(checker.ConfigDir)
|
|
if err != nil {
|
|
return
|
|
}
|
|
now := checker.now()
|
|
for _, node := range store.Nodes {
|
|
if !IsHealthCheckDue(node, now) {
|
|
continue
|
|
}
|
|
nodeID := node.ID
|
|
go func() {
|
|
_, _ = checker.CheckNodeByID(nodeID, "scheduler")
|
|
}()
|
|
}
|
|
}
|
|
|
|
// IsHealthCheckDue reports whether a node should be probed now.
|
|
func IsHealthCheckDue(node settings.Node, now time.Time) bool {
|
|
if node.HealthCheckIntervalSeconds <= 0 {
|
|
return false
|
|
}
|
|
if node.HealthLastCheckedAt == nil {
|
|
return true
|
|
}
|
|
elapsed := now.Sub(node.HealthLastCheckedAt.UTC())
|
|
return elapsed >= time.Duration(node.HealthCheckIntervalSeconds)*time.Second
|
|
}
|
|
|
|
// CheckNodeByID loads credentials, runs ping + SSH, and persists the result.
|
|
func (checker *Checker) CheckNodeByID(nodeID string, actor string) (settings.Node, error) {
|
|
if !checker.tryLockNode(nodeID) {
|
|
store, err := settings.LoadNodesOrEmpty(checker.ConfigDir)
|
|
if err != nil {
|
|
return settings.Node{}, err
|
|
}
|
|
for _, node := range store.Nodes {
|
|
if node.ID == nodeID {
|
|
return node, fmt.Errorf("health check already running for node")
|
|
}
|
|
}
|
|
return settings.Node{}, fmt.Errorf("node not found")
|
|
}
|
|
defer checker.unlockNode(nodeID)
|
|
|
|
if len(checker.Key) == 0 {
|
|
return settings.Node{}, fmt.Errorf("%s is not set", settings.ConfigKeyEnvVar)
|
|
}
|
|
|
|
store, err := settings.LoadNodesOrEmpty(checker.ConfigDir)
|
|
if err != nil {
|
|
return settings.Node{}, err
|
|
}
|
|
|
|
nodeIndex := -1
|
|
var node settings.Node
|
|
for index, candidate := range store.Nodes {
|
|
if candidate.ID != nodeID {
|
|
continue
|
|
}
|
|
node = candidate
|
|
nodeIndex = index
|
|
break
|
|
}
|
|
if nodeIndex < 0 {
|
|
return settings.Node{}, fmt.Errorf("node not found")
|
|
}
|
|
|
|
keyStore, err := settings.LoadNodeKeysOrEmpty(checker.ConfigDir, checker.Key)
|
|
if err != nil {
|
|
return settings.Node{}, err
|
|
}
|
|
|
|
var keyEntry settings.NodeKeyEntry
|
|
keyFound := false
|
|
for _, entry := range keyStore.Keys {
|
|
if entry.NodeID == node.ID {
|
|
keyEntry = entry
|
|
keyFound = true
|
|
break
|
|
}
|
|
}
|
|
if !keyFound {
|
|
return settings.Node{}, fmt.Errorf("private key not found for node")
|
|
}
|
|
|
|
result := checker.runChecks(node, keyEntry)
|
|
updated := applyHealthResult(node, result)
|
|
store.Nodes[nodeIndex] = updated
|
|
if err := settings.SaveNodes(checker.ConfigDir, store); err != nil {
|
|
return settings.Node{}, err
|
|
}
|
|
|
|
_ = appendHealthAudit(checker.ConfigDir, actor, updated, result.Message)
|
|
return updated, nil
|
|
}
|
|
|
|
func (checker *Checker) runChecks(node settings.Node, keyEntry settings.NodeKeyEntry) Result {
|
|
checkedAt := checker.now()
|
|
pingFn := checker.PingFn
|
|
if pingFn == nil {
|
|
pingFn = auth.PingHost
|
|
}
|
|
sshFn := checker.SSHFn
|
|
if sshFn == nil {
|
|
sshFn = auth.TestSSHConnection
|
|
}
|
|
|
|
pingResult := pingFn(node.HostIP)
|
|
sshErr := sshFn(node.HostIP, node.Username, keyEntry.PrivateKey, keyEntry.Passphrase)
|
|
sshOK := sshErr == nil
|
|
|
|
message := buildHealthMessage(pingResult, sshOK, sshErr)
|
|
return Result{
|
|
PingOK: pingResult.OK,
|
|
SSHOK: sshOK,
|
|
OK: pingResult.OK && sshOK,
|
|
Message: message,
|
|
CheckedAt: checkedAt,
|
|
}
|
|
}
|
|
|
|
func buildHealthMessage(pingResult auth.PingResult, sshOK bool, sshErr error) string {
|
|
var parts []string
|
|
if pingResult.OK {
|
|
parts = append(parts, fmt.Sprintf("ping ok (%s)", pingResult.Method))
|
|
} else if pingResult.Error != nil {
|
|
parts = append(parts, fmt.Sprintf("ping failed: %v", pingResult.Error))
|
|
} else {
|
|
parts = append(parts, "ping failed")
|
|
}
|
|
|
|
if sshOK {
|
|
parts = append(parts, "ssh ok")
|
|
} else if sshErr != nil {
|
|
parts = append(parts, fmt.Sprintf("ssh failed: %v", sshErr))
|
|
} else {
|
|
parts = append(parts, "ssh failed")
|
|
}
|
|
|
|
return fmt.Sprintf("%s; %s", parts[0], parts[1])
|
|
}
|
|
|
|
func applyHealthResult(node settings.Node, result Result) settings.Node {
|
|
checkedAt := result.CheckedAt
|
|
pingOK := result.PingOK
|
|
sshOK := result.SSHOK
|
|
healthOK := result.OK
|
|
node.HealthLastCheckedAt = &checkedAt
|
|
node.HealthPingOK = &pingOK
|
|
node.HealthSSHOK = &sshOK
|
|
node.HealthOK = &healthOK
|
|
node.HealthMessage = result.Message
|
|
return node
|
|
}
|
|
|
|
func appendHealthAudit(configDir string, actor string, node settings.Node, detail string) error {
|
|
eventID, err := auth.NewUUID()
|
|
if err != nil {
|
|
eventID = fmt.Sprintf("health-%d", time.Now().UnixNano())
|
|
}
|
|
return settings.AppendNodeAuditEvent(configDir, settings.NodeAuditEvent{
|
|
ID: eventID,
|
|
At: time.Now().UTC(),
|
|
Action: settings.NodeAuditActionHealthCheck,
|
|
Actor: actor,
|
|
NodeID: node.ID,
|
|
NodeName: node.Name,
|
|
NodeKind: node.Kind,
|
|
Detail: detail,
|
|
})
|
|
}
|