Files
ClusterCanvas/service/internal/runner/executor.go
T
Squid 0e06063c1d Clarify run-if skip reasons with the resolved variable value.
Failed conditions now explain skips in plain language (e.g. UpdateCount was 0) and truncate long captures so misconfigured apt listings cannot flood the UI.
2026-07-19 21:59:35 +02:00

450 lines
12 KiB
Go

package runner
import (
"fmt"
"strings"
"sync"
"time"
"codeberg.org/SquidSE/ClusterCanvas/service/internal/auth"
"codeberg.org/SquidSE/ClusterCanvas/service/internal/settings"
)
const defaultRunTimeout = 5 * time.Minute
// Executor runs action groups and items against nodes over SSH.
type Executor struct {
ConfigDir string
Key []byte
mu sync.Mutex
running map[string]struct{} // groupID -> running
stopCh chan struct{}
stopped chan struct{}
tickerOn bool
}
// NewExecutor creates an Executor for the given config directory and encryption key.
func NewExecutor(configDir string, key []byte) *Executor {
return &Executor{
ConfigDir: configDir,
Key: key,
running: map[string]struct{}{},
stopCh: make(chan struct{}),
stopped: make(chan struct{}),
}
}
// StartScheduler begins the background schedule ticker (30s).
func (executor *Executor) StartScheduler() {
executor.mu.Lock()
if executor.tickerOn {
executor.mu.Unlock()
return
}
executor.tickerOn = true
executor.mu.Unlock()
go func() {
defer close(executor.stopped)
ticker := time.NewTicker(30 * time.Second)
defer ticker.Stop()
for {
select {
case <-executor.stopCh:
return
case <-ticker.C:
executor.tickSchedules()
}
}
}()
}
// StopScheduler stops the background ticker and waits for it to exit.
func (executor *Executor) StopScheduler() {
executor.mu.Lock()
if !executor.tickerOn {
executor.mu.Unlock()
return
}
executor.tickerOn = false
executor.mu.Unlock()
close(executor.stopCh)
<-executor.stopped
}
func (executor *Executor) tryLockGroup(groupID string) bool {
executor.mu.Lock()
defer executor.mu.Unlock()
if _, ok := executor.running[groupID]; ok {
return false
}
executor.running[groupID] = struct{}{}
return true
}
func (executor *Executor) unlockGroup(groupID string) {
executor.mu.Lock()
defer executor.mu.Unlock()
delete(executor.running, groupID)
}
func (executor *Executor) tickSchedules() {
store, err := settings.LoadNodeActionGroupsOrEmpty(executor.ConfigDir)
if err != nil {
return
}
now := time.Now()
for _, group := range store.Groups {
if !settings.IsScheduleDue(group, now) {
continue
}
groupCopy := group
go func() {
_, _ = executor.RunGroup(groupCopy, settings.ActionRunTriggerSchedule, "scheduler")
}()
}
}
// ResolvedAction is the effective action definition after library resolution.
type ResolvedAction struct {
Name string
Description string
Kind settings.ActionKind
Body string
Env []settings.ActionEnvVar
RequiresSudo bool
Source settings.ActionItemSource
}
// ResolveItem resolves a library or local item to an executable action.
func ResolveItem(item settings.NodeActionItem, library []settings.Action) (ResolvedAction, error) {
switch item.Source {
case settings.ActionItemSourceLibrary:
if strings.TrimSpace(item.LibraryActionID) == "" {
return ResolvedAction{}, fmt.Errorf("library action id is required")
}
for _, action := range library {
if action.ID == item.LibraryActionID {
env := action.Env
if env == nil {
env = []settings.ActionEnvVar{}
}
return ResolvedAction{
Name: action.Name,
Description: action.Description,
Kind: action.Kind,
Body: action.Body,
Env: env,
RequiresSudo: action.RequiresSudo,
Source: settings.ActionItemSourceLibrary,
}, nil
}
}
return ResolvedAction{}, fmt.Errorf("library action %q not found", item.LibraryActionID)
case settings.ActionItemSourceLocal:
if strings.TrimSpace(item.Name) == "" {
return ResolvedAction{}, fmt.Errorf("local action name is required")
}
if strings.TrimSpace(item.Body) == "" {
return ResolvedAction{}, fmt.Errorf("local action body is required")
}
if item.Kind != settings.ActionKindShell && item.Kind != settings.ActionKindScript {
return ResolvedAction{}, fmt.Errorf("local action kind must be shell or script")
}
env := item.Env
if env == nil {
env = []settings.ActionEnvVar{}
}
return ResolvedAction{
Name: item.Name,
Description: item.Description,
Kind: item.Kind,
Body: item.Body,
Env: env,
RequiresSudo: item.RequiresSudo,
Source: settings.ActionItemSourceLocal,
}, nil
default:
return ResolvedAction{}, fmt.Errorf("unknown action item source %q", item.Source)
}
}
// WrapShellWithSudo prepends "sudo -n " when requiresSudo is true.
func WrapShellWithSudo(command string, requiresSudo bool) string {
if !requiresSudo {
return command
}
return "sudo -n " + command
}
// ScriptRemoteCommand is the remote argv used for script actions (body on stdin).
func ScriptRemoteCommand(requiresSudo bool) string {
if requiresSudo {
return "sudo -n bash -s"
}
return "bash -s"
}
// BuildRemoteCommand returns the final SSH exec string for an action (including env exports).
func BuildRemoteCommand(
kind settings.ActionKind,
body string,
requiresSudo bool,
env map[string]string,
) string {
var command string
switch kind {
case settings.ActionKindShell:
command = WrapShellWithSudo(body, requiresSudo)
case settings.ActionKindScript:
command = ScriptRemoteCommand(requiresSudo)
default:
command = body
}
return auth.PrependEnvExports(env, command)
}
// RunGroup executes all items in order, appends a log file entry, and updates LastRunAt.
func (executor *Executor) RunGroup(
group settings.NodeActionGroup,
trigger settings.ActionRunTrigger,
actor string,
) (settings.ActionGroupRunRecord, error) {
if !executor.tryLockGroup(group.ID) {
return settings.ActionGroupRunRecord{}, fmt.Errorf("action group is already running")
}
defer executor.unlockGroup(group.ID)
return executor.runItems(group, group.Items, trigger, actor)
}
// RunItem executes a single item (still logged under the parent group file).
func (executor *Executor) RunItem(
group settings.NodeActionGroup,
item settings.NodeActionItem,
trigger settings.ActionRunTrigger,
actor string,
) (settings.ActionGroupRunRecord, error) {
lockKey := group.ID + ":" + item.ID
if !executor.tryLockGroup(lockKey) {
return settings.ActionGroupRunRecord{}, fmt.Errorf("action is already running")
}
defer executor.unlockGroup(lockKey)
return executor.runItems(group, []settings.NodeActionItem{item}, trigger, actor)
}
func (executor *Executor) runItems(
group settings.NodeActionGroup,
items []settings.NodeActionItem,
trigger settings.ActionRunTrigger,
actor string,
) (settings.ActionGroupRunRecord, error) {
runID, err := auth.NewUUID()
if err != nil {
return settings.ActionGroupRunRecord{}, err
}
node, privateKey, passphrase, err := executor.loadNodeCredentials(group.NodeID)
if err != nil {
return settings.ActionGroupRunRecord{}, err
}
libraryStore, err := settings.LoadActionsOrSeed(executor.ConfigDir)
if err != nil {
return settings.ActionGroupRunRecord{}, err
}
startedAt := time.Now().UTC()
results := make([]settings.ActionItemRunResult, 0, len(items))
runVars := map[string]string{}
for _, item := range items {
itemStarted := time.Now().UTC()
resolved, resolveErr := ResolveItem(item, libraryStore.Actions)
if resolveErr != nil {
results = append(results, settings.ActionItemRunResult{
ItemID: item.ID,
ActionName: item.Name,
Source: string(item.Source),
ExitCode: -1,
Stdout: "",
Stderr: "",
Error: resolveErr.Error(),
StartedAt: itemStarted,
FinishedAt: time.Now().UTC(),
})
continue
}
runIf := strings.TrimSpace(item.RunIf)
if runIf != "" {
shouldRun, detail, condErr := EvaluateCondition(runIf, runVars)
if condErr != nil || !shouldRun {
skipReason := fmt.Sprintf("%s skipped: run condition was not met", resolved.Name)
if detail != "" {
skipReason = fmt.Sprintf("%s skipped: %s", resolved.Name, detail)
}
if condErr != nil {
skipReason = fmt.Sprintf("%s skipped because condition error: %s", resolved.Name, condErr.Error())
}
results = append(results, settings.ActionItemRunResult{
ItemID: item.ID,
ActionName: resolved.Name,
Source: string(resolved.Source),
ExitCode: 0,
Stdout: "",
Stderr: "",
Skipped: true,
SkipReason: skipReason,
StartedAt: itemStarted,
FinishedAt: time.Now().UTC(),
})
continue
}
}
rawEnv := make(map[string]string, len(resolved.Env))
for _, envVar := range resolved.Env {
rawEnv[envVar.Name] = envVar.Value
}
envMap := make(map[string]string, len(resolved.Env))
for _, envVar := range resolved.Env {
expandedValue := auth.ExpandPlaceholders(envVar.Value, auth.PlaceholderContext{
Host: node.Name,
IP: node.HostIP,
Username: node.Username,
Env: rawEnv,
Secrets: map[string]string{},
})
envMap[envVar.Name] = ExpandDollarVars(expandedValue, runVars)
}
expandedBody := auth.ExpandPlaceholders(resolved.Body, auth.PlaceholderContext{
Host: node.Name,
IP: node.HostIP,
Username: node.Username,
Env: envMap,
Secrets: map[string]string{},
})
expandedBody = ExpandDollarVars(expandedBody, runVars)
remoteCommand := BuildRemoteCommand(
resolved.Kind,
expandedBody,
resolved.RequiresSudo,
envMap,
)
var sshResult auth.SSHRunResult
var runErr error
switch resolved.Kind {
case settings.ActionKindShell:
sshResult, runErr = auth.RunSSHShell(
node.HostIP,
node.Username,
privateKey,
passphrase,
WrapShellWithSudo(expandedBody, resolved.RequiresSudo),
envMap,
defaultRunTimeout,
)
case settings.ActionKindScript:
sshResult, runErr = auth.RunSSHScript(
node.HostIP,
node.Username,
privateKey,
passphrase,
expandedBody,
envMap,
defaultRunTimeout,
resolved.RequiresSudo,
)
default:
runErr = fmt.Errorf("unsupported action kind %q", resolved.Kind)
}
itemResult := settings.ActionItemRunResult{
ItemID: item.ID,
ActionName: resolved.Name,
Source: string(resolved.Source),
Kind: resolved.Kind,
RemoteCommand: remoteCommand,
SSHUsername: node.Username,
ExitCode: sshResult.ExitCode,
Stdout: sshResult.Stdout,
Stderr: sshResult.Stderr,
StartedAt: itemStarted,
FinishedAt: time.Now().UTC(),
}
if runErr != nil {
itemResult.Error = runErr.Error()
if itemResult.ExitCode == 0 {
itemResult.ExitCode = -1
}
} else if setName := strings.TrimSpace(item.SetVariable); setName != "" {
captured := strings.TrimSpace(sshResult.Stdout)
runVars[setName] = captured
itemResult.SetVariable = setName
itemResult.VariableValue = captured
}
results = append(results, itemResult)
}
finishedAt := time.Now().UTC()
record := settings.ActionGroupRunRecord{
ID: runID,
NodeID: group.NodeID,
GroupID: group.ID,
GroupName: group.Name,
Trigger: trigger,
Actor: actor,
StartedAt: startedAt,
FinishedAt: finishedAt,
Actions: results,
}
if err := settings.AppendActionLogRun(executor.ConfigDir, record); err != nil {
return record, fmt.Errorf("append action log: %w", err)
}
if err := settings.UpdateGroupLastRunAt(executor.ConfigDir, group.ID, finishedAt); err != nil {
return record, fmt.Errorf("update last run: %w", err)
}
return record, nil
}
func (executor *Executor) loadNodeCredentials(nodeID string) (settings.Node, string, string, error) {
if len(executor.Key) == 0 {
return settings.Node{}, "", "", fmt.Errorf("%s is not set", settings.ConfigKeyEnvVar)
}
nodeStore, err := settings.LoadNodesOrEmpty(executor.ConfigDir)
if err != nil {
return settings.Node{}, "", "", err
}
var node settings.Node
found := false
for _, candidate := range nodeStore.Nodes {
if candidate.ID == nodeID {
node = candidate
found = true
break
}
}
if !found {
return settings.Node{}, "", "", fmt.Errorf("node not found")
}
keyStore, err := settings.LoadNodeKeysOrEmpty(executor.ConfigDir, executor.Key)
if err != nil {
return settings.Node{}, "", "", err
}
for _, entry := range keyStore.Keys {
if entry.NodeID == nodeID {
return node, entry.PrivateKey, entry.Passphrase, nil
}
}
return settings.Node{}, "", "", fmt.Errorf("node private key not found")
}