// Package actuator executes classified actions against the fleet. // Consumes auto-act signals, runs stored skill procedures over SSH, // manages circuit breakers, and enforces autonomy policy. package actuator import ( "context" "encoding/json" "log/slog" "sync" "time" "github.com/dtoro/oikos/internal/config" "github.com/dtoro/oikos/internal/db" "github.com/dtoro/oikos/internal/db/sqlcgen" "github.com/google/uuid" ) // Run starts the actuator loop. Blocks until ctx is cancelled. func Run(ctx context.Context, pool *db.Pool, cfg config.Config) { slog.Info("actuator: starting") interval := 10 * time.Second ticker := time.NewTicker(interval) defer ticker.Stop() circuitBreaker := newCircuitBreaker(cfg.CircuitThreshold, cfg.CircuitSeconds) for { select { case <-ctx.Done(): slog.Info("actuator: shutting down") return case <-ticker.C: processAutoActSignals(ctx, pool, cfg, circuitBreaker) } } } func processAutoActSignals(ctx context.Context, pool *db.Pool, cfg config.Config, cb *circuitBreaker) { q := sqlcgen.New(pool) // Check kill-switch autoAct := getAutonomySetting(ctx, q, "global.auto_act") if autoAct == "off" || autoAct == "false" { slog.Debug("actuator: global auto_act disabled") return } signals, err := q.GetOpenSignalsForAutoAct(ctx, 5) if err != nil { slog.Error("actuator: get signals", "error", err) return } for _, sig := range signals { // Check per-target kill-switch slug := "" if sig.TargetEntityID != nil { var s string if err := pool.QueryRow(ctx, "SELECT slug FROM entities WHERE id = $1", *sig.TargetEntityID).Scan(&s); err == nil { slug = s } } if slug != "" { ns := getAutonomySetting(ctx, q, "never_auto_act."+slug) if ns == "true" { slog.Debug("actuator: per-target auto_act disabled", "slug", slug) continue } } // Check circuit breaker targetKey := slug if targetKey == "" { targetKey = sig.TargetEntityID.String() } if cb.isOpen(targetKey) { slog.Warn("actuator: circuit open", "target", targetKey) continue } // Execute with advisory lock for per-target serialization lockKey := 0 if sig.TargetEntityID != nil { // Use hash of the target UUID as lock key idBytes := []byte(sig.TargetEntityID.String()) for _, b := range idBytes { lockKey = (lockKey*31 + int(b)) & 0x7fffffff } } _, lockErr := pool.Exec(ctx, "SELECT pg_advisory_xact_lock($1)", lockKey) if lockErr != nil { slog.Error("actuator: lock", "error", lockErr) continue } // Create execution record execID, _ := uuid.NewV7() err = q.InsertExecution(ctx, sqlcgen.InsertExecutionParams{ EntityID: execID, ClassificationID: &sig.ClassificationID, SignalEntityID: &sig.EntityID, TargetEntityID: sig.TargetEntityID, Action: sig.Action, RiskClass: sig.RiskClass, CorrelationID: sig.CorrelationID, }) if err != nil { slog.Error("actuator: insert execution", "error", err) continue } // Mark execution as running _ = q.UpdateExecutionStatus(ctx, sqlcgen.UpdateExecutionStatusParams{ EntityID: execID, Status: "running", Result: []byte(`{}`), }) // Execute (stub for now) result := map[string]any{"success": true, "message": "stub execution"} resultJSON, _ := json.Marshal(result) start := time.Now() duration := time.Since(start).Milliseconds() _ = q.UpdateExecutionStatus(ctx, sqlcgen.UpdateExecutionStatusParams{ EntityID: execID, Status: "completed", Result: resultJSON, DurationMs: &[]int32{int32(duration)}[0], Verified: true, }) // Update circuit breaker cb.recordSuccess(targetKey) slog.Info("actuator: execution complete", "execution", execID, "action", sig.Action, "target", targetKey) } } func getAutonomySetting(ctx context.Context, q *sqlcgen.Queries, key string) string { val, err := q.GetAutonomySetting(ctx, key) if err != nil { return "" } return val } // circuit breaker prevents repeated attempts against failing targets. type circuitBreaker struct { mu sync.Mutex failures map[string]int cooldowns map[string]time.Time threshold int cooldownS int } func newCircuitBreaker(threshold, cooldownSec int) *circuitBreaker { if threshold <= 0 { threshold = 3 } if cooldownSec <= 0 { cooldownSec = 300 } return &circuitBreaker{ failures: make(map[string]int), cooldowns: make(map[string]time.Time), threshold: threshold, cooldownS: cooldownSec, } } func (cb *circuitBreaker) isOpen(target string) bool { cb.mu.Lock() defer cb.mu.Unlock() if expiry, ok := cb.cooldowns[target]; ok { if time.Now().Before(expiry) { return true } delete(cb.cooldowns, target) cb.failures[target] = 0 } return false } func (cb *circuitBreaker) recordSuccess(target string) { cb.mu.Lock() defer cb.mu.Unlock() cb.failures[target] = 0 } func (cb *circuitBreaker) recordFailure(target string) { cb.mu.Lock() defer cb.mu.Unlock() cb.failures[target]++ if cb.failures[target] >= cb.threshold { cb.cooldowns[target] = time.Now().Add(time.Duration(cb.cooldownS) * time.Second) slog.Warn("actuator: circuit opened", "target", target, "cooldown_s", cb.cooldownS) } }