goherence/internal/executor/executor.go
2026-09-11 10:17:25 +03:00

472 lines
17 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// Package executor связывает все остальные пакеты воедино: парсер даёт
// структуру плейбука, lint проверяет её на guardrails, vars считает
// переменные для каждого хоста, module.Registry резолвит модуль по имени,
// ssh.Conn исполняет его на хосте. Сам executor не содержит модульной
// логики — только оркестрацию.
package executor
import (
"context"
"fmt"
"strings"
"sync"
"github.com/vladimir/goherence/internal/facts"
"github.com/vladimir/goherence/internal/lint"
"github.com/vladimir/goherence/internal/module"
"github.com/vladimir/goherence/internal/module/builtin"
"github.com/vladimir/goherence/internal/module/external"
"github.com/vladimir/goherence/internal/parser"
"github.com/vladimir/goherence/internal/ssh"
tmpl "github.com/vladimir/goherence/internal/template"
"github.com/vladimir/goherence/internal/vars"
)
// Options — то, что задаётся снаружи (флагами CLI/конфигом) при запуске плейбука.
type Options struct {
Limit string // ограничить выполнение одним хостом/группой
ExtraVars map[string]interface{}
CheckMode bool
Force bool // пропустить блокирующие ошибки линтера
ModulesDir string // путь к modules.d/ с внешними модулями; пусто — внешних модулей нет
Forks int // сколько хостов внутри одного батча обрабатывать одновременно; 0/1 — последовательно
Serial int // размер батча rolling-обновления; 0 — один батч на всех хостов сразу
MaxFails int // после скольких упавших хостов прервать оставшиеся батчи; 0 — без ограничения
Become bool // глобально поднимать привилегии (sudo) для всех команд на всех хостах
GatherFacts bool // собирать ли internal/facts перед тасками
OnTaskStart func(host string, t *parser.Task)
OnTaskEnd func(host string, t *parser.Task, res module.Result)
}
// SSHConfigFor — функция, которую вызывающий код должен предоставить,
// чтобы получить параметры подключения для конкретного хоста (ключ,
// пользователь и т.д.). Executor намеренно не знает, откуда эти данные
// берутся (инвентарь / vault / переменные окружения) — это ответственность
// вызывающего кода (см. cmd/goherence).
type SSHConfigFor func(host string) ssh.Config
// Run выполняет плейбук pb против инвентаря inv. Хосты разбиваются на
// батчи по opts.Serial (rolling-обновление — "один хост за раз" это
// Serial: 1); внутри батча до opts.Forks хостов обрабатываются
// параллельно. После каждого батча проверяется opts.MaxFails: если
// упавших хостов уже достаточно, оставшиеся батчи не запускаются —
// то, что уже успело обновиться, не откатывается (это не транзакция).
func Run(ctx context.Context, pb *parser.Playbook, inv *parser.Inventory, sshCfg SSHConfigFor, opts Options) error {
violations := lint.New().Run(pb)
for _, v := range violations {
fmt.Printf("[%s] %s: %s: %s\n", v.Severity, v.Rule, v.Location, v.Message)
}
if lint.HasError(violations) && !opts.Force {
return fmt.Errorf("плейбук нарушает guardrails декларативности; " +
"используй --force для запуска несмотря на ошибки (не рекомендуется)")
}
// printMu сериализует вывод OnTaskStart/OnTaskEnd между горутинами —
// без этого строки от разных хостов перемежались бы посимвольно.
printMu := &sync.Mutex{}
origStart, origEnd := opts.OnTaskStart, opts.OnTaskEnd
if origStart != nil {
opts.OnTaskStart = func(host string, t *parser.Task) {
printMu.Lock()
defer printMu.Unlock()
origStart(host, t)
}
}
if origEnd != nil {
opts.OnTaskEnd = func(host string, t *parser.Task, res module.Result) {
printMu.Lock()
defer printMu.Unlock()
origEnd(host, t, res)
}
}
for _, play := range pb.Plays {
hosts := resolveHosts(inv, play.Hosts, opts.Limit)
if err := runPlaySerially(ctx, pb, inv, &play, hosts, sshCfg, opts); err != nil {
return err
}
}
return nil
}
// runPlaySerially делит hosts на батчи по opts.Serial и прогоняет их
// последовательно один за другим, проверяя opts.MaxFails между батчами.
func runPlaySerially(
ctx context.Context,
pb *parser.Playbook,
inv *parser.Inventory,
play *parser.Play,
hosts []string,
sshCfg SSHConfigFor,
opts Options,
) error {
batchSize := opts.Serial
if batchSize < 1 {
batchSize = len(hosts) // без serial — один батч на всех, как раньше
}
var failed []error
for start := 0; start < len(hosts); start += batchSize {
end := start + batchSize
if end > len(hosts) {
end = len(hosts)
}
batch := hosts[start:end]
failed = append(failed, runBatch(ctx, pb, inv, play, batch, sshCfg, opts)...)
if opts.MaxFails > 0 && len(failed) >= opts.MaxFails {
return fmt.Errorf(
"остановлено после %d ошибок (лимит max_fails=%d), оставшиеся хосты не тронуты: %s",
len(failed), opts.MaxFails, joinErrors(failed))
}
}
if len(failed) > 0 {
return fmt.Errorf("ошибки на %d из %d хостов: %s", len(failed), len(hosts), joinErrors(failed))
}
return nil
}
// runBatch запускает play на hosts, максимум opts.Forks одновременно,
// и возвращает ВСЕ ошибки батча (а не первую попавшуюся) — иначе после
// каждого батча нечего было бы сравнивать с MaxFails: одна необработанная
// ошибка ничем не отличалась бы от трёх.
func runBatch(
ctx context.Context,
pb *parser.Playbook,
inv *parser.Inventory,
play *parser.Play,
hosts []string,
sshCfg SSHConfigFor,
opts Options,
) []error {
forks := opts.Forks
if forks < 1 {
forks = 1
}
if forks == 1 {
var errs []error
for _, host := range hosts {
if err := runPlayOnHost(ctx, pb, inv, play, host, sshCfg, opts); err != nil {
errs = append(errs, fmt.Errorf("хост %s: %w", host, err))
}
}
return errs
}
sem := make(chan struct{}, forks)
errCh := make(chan error, len(hosts))
var wg sync.WaitGroup
for _, host := range hosts {
host := host
wg.Add(1)
sem <- struct{}{}
go func() {
defer wg.Done()
defer func() { <-sem }()
if err := runPlayOnHost(ctx, pb, inv, play, host, sshCfg, opts); err != nil {
errCh <- fmt.Errorf("хост %s: %w", host, err)
}
}()
}
wg.Wait()
close(errCh)
var errs []error
for err := range errCh {
errs = append(errs, err)
}
return errs
}
func joinErrors(errs []error) string {
parts := make([]string, len(errs))
for i, e := range errs {
parts[i] = e.Error()
}
return strings.Join(parts, "; ")
}
// resolveHosts разворачивает play.Hosts (имя группы или "all") в список
// конкретных хостов, применяя --limit, если он задан.
func resolveHosts(inv *parser.Inventory, hostsSpec, limit string) []string {
group, ok := inv.Groups[hostsSpec]
if !ok {
return nil
}
var out []string
for name := range group.Hosts {
if limit == "" || limit == name || limit == hostsSpec {
out = append(out, name)
}
}
return out
}
func runPlayOnHost(
ctx context.Context,
pb *parser.Playbook,
inv *parser.Inventory,
play *parser.Play,
host string,
sshCfg SSHConfigFor,
opts Options,
) error {
conn, err := ssh.Dial(sshCfg(host))
if err != nil {
return fmt.Errorf("подключение по ssh: %w", err)
}
defer conn.Close()
conn.Become = opts.Become // глобальный переключатель на весь прогон, см. internal/config
// факты собираются один раз на хост в начале play — ровно как
// implicit gather_facts в Ansible перед первым таском. Отключается
// через opts.GatherFacts=false (config: gather_facts: false), если
// плейбуку факты не нужны — это экономит одну SSH-сессию на хост.
hostFacts := map[string]interface{}{}
if opts.GatherFacts {
hostFacts, err = facts.Gather(conn)
if err != nil {
return fmt.Errorf("сбор фактов: %w", err)
}
}
// notified — хэндлеры, на которые сработал notify хотя бы одного
// изменившего состояние таска, в порядке первого срабатывания.
// Собирается по всему play (play-level tasks + все роли) и
// выполняется один раз в конце — ровно так же, как в Ansible.
var notified []string
seen := map[string]bool{}
notify := func(names []string) {
for _, n := range names {
if !seen[n] {
seen[n] = true
notified = append(notified, n)
}
}
}
// сначала таски самого play, затем таски каждой подключённой роли —
// в первой версии roles всегда выполняются после play-level tasks
n, err := runTasks(ctx, play.Tasks, nil, inv, play, host, conn, hostFacts, opts)
if err != nil {
return err
}
notify(n)
handlers := map[string]parser.Task{}
for _, roleName := range play.Roles {
role := pb.Roles[roleName]
if role == nil {
return fmt.Errorf("роль %q не загружена", roleName)
}
for _, h := range role.Handlers {
handlers[h.Name] = h
}
n, err := runTasks(ctx, role.Tasks, role, inv, play, host, conn, hostFacts, opts)
if err != nil {
return fmt.Errorf("роль %s: %w", roleName, err)
}
notify(n)
}
return runHandlers(ctx, notified, handlers, inv, play, host, conn, hostFacts, opts)
}
// runHandlers выполняет по одному разу каждый хэндлер из notified —
// аналог того, как Ansible прогоняет notify-хэндлеры в конце play.
// Хэндлер резолвится с переменными play/host/facts, но без role.Defaults/Vars
// конкретной роли (в первой версии хэндлеры общие на весь play, без
// привязки к переменным той роли, где они были объявлены) — сознательное
// упрощение, отмеченное в README.
func runHandlers(
ctx context.Context,
notified []string,
handlers map[string]parser.Task,
inv *parser.Inventory,
play *parser.Play,
host string,
conn *ssh.Conn,
hostFacts map[string]interface{},
opts Options,
) error {
if len(notified) == 0 {
return nil
}
var toRun []parser.Task
for _, name := range notified {
h, ok := handlers[name]
if !ok {
return fmt.Errorf("notify указывает на неизвестный хэндлер %q", name)
}
toRun = append(toRun, h)
}
_, err := runTasks(ctx, toRun, nil, inv, play, host, conn, hostFacts, opts)
return err
}
// runTasks выполняет список тасков на одном хосте и возвращает имена
// хэндлеров, на которые сработал notify (для тасков, реально изменивших
// состояние — Result.Changed). Вызывающий код (runPlayOnHost) собирает
// notify со всех групп тасков play и роли и прогоняет хэндлеры один раз в конце.
func runTasks(
ctx context.Context,
tasks []parser.Task,
role *parser.Role,
inv *parser.Inventory,
play *parser.Play,
host string,
conn *ssh.Conn,
hostFacts map[string]interface{},
opts Options,
) ([]string, error) {
hostVars := vars.Resolve(host, inv, play, role, hostFacts, opts.ExtraVars)
// граф зависимостей (after/before) — см. graph.go; без объявленных
// зависимостей порядок не меняется вообще.
orderedTasks, err := orderTasks(tasks)
if err != nil {
return nil, err
}
tasks = orderedTasks
reg := module.NewRegistry()
builtin.RegisterAll(reg, hostVars) // TemplateModule/PackageModule получают vars именно здесь,
// на пересчёт для каждого набора тасков — см. ограничение в register.go
if opts.ModulesDir != "" {
extModules, err := external.LoadDir(opts.ModulesDir)
if err != nil {
return nil, fmt.Errorf("загрузка внешних модулей из %s: %w", opts.ModulesDir, err)
}
for _, m := range extModules {
reg.RegisterExternal(m)
}
}
var notified []string
for i := range tasks {
t := &tasks[i]
if t.When != "" && !evalWhen(t.When, hostVars) {
continue
}
mod, err := reg.Resolve(t.Module)
if err != nil {
return nil, err
}
items, err := loopItems(t.Loop, hostVars)
if err != nil {
return nil, fmt.Errorf("таск %q: loop: %w", t.Name, err)
}
for _, item := range items {
args, err := renderTaskArgs(t.Args, hostVars, item)
if err != nil {
return nil, fmt.Errorf("таск %q: рендер аргументов: %w", t.Name, err)
}
if opts.OnTaskStart != nil {
opts.OnTaskStart(host, t)
}
res, err := mod.Run(ctx, module.Input{
SchemaVersion: module.SchemaVersion,
Args: args,
Host: host,
CheckMode: opts.CheckMode,
}, conn)
if err != nil {
return nil, fmt.Errorf("таск %q: %w", t.Name, err)
}
if opts.OnTaskEnd != nil {
opts.OnTaskEnd(host, t, res)
}
if res.Failed {
return nil, fmt.Errorf("таск %q завершился с ошибкой: %s", t.Name, res.Msg)
}
if res.Changed && len(t.Notify) > 0 {
notified = append(notified, t.Notify...)
}
}
}
return notified, nil
}
// loopItems разворачивает t.Loop в список элементов, по одному на итерацию.
// nil-loop — это ровно один "элемент" (пустой), то есть таск выполняется
// один раз как обычно; это позволяет не разветвлять код на "с loop"/"без loop".
func loopItems(loopSpec interface{}, hostVars map[string]interface{}) ([]interface{}, error) {
if loopSpec == nil {
return []interface{}{nil}, nil
}
switch v := loopSpec.(type) {
case []interface{}:
return v, nil
case string:
// "{{ .some_list }}" — ссылка на переменную-список в hostVars
varName := extractVarName(v)
if varName == "" {
return nil, fmt.Errorf("loop-выражение %q не распознано (первая версия понимает только {{ .var }})", v)
}
val, ok := hostVars[varName]
if !ok {
return nil, fmt.Errorf("переменная %q для loop не найдена", varName)
}
list, ok := val.([]interface{})
if !ok {
return nil, fmt.Errorf("переменная %q не является списком", varName)
}
return list, nil
default:
return nil, fmt.Errorf("неподдерживаемый тип loop: %T", loopSpec)
}
}
// extractVarName вытаскивает имя переменной из "{{ .name }}" — только
// этот один паттерн, без произвольных выражений (loop — не место для логики).
func extractVarName(expr string) string {
s := strings.TrimSpace(expr)
s = strings.TrimPrefix(s, "{{")
s = strings.TrimSuffix(s, "}}")
s = strings.TrimSpace(s)
s = strings.TrimPrefix(s, ".")
return strings.TrimSpace(s)
}
// renderTaskArgs рендерит каждое строковое значение в args через
// internal/template (text/template+sprig), добавляя текущий item из
// loop под именем `.item` — ровно так же, как Ansible подставляет `item`
// внутри тасков с loop. Нестроковые значения (числа, bool, вложенные map)
// передаются как есть.
func renderTaskArgs(args map[string]interface{}, hostVars map[string]interface{}, item interface{}) (map[string]interface{}, error) {
scopedVars := make(map[string]interface{}, len(hostVars)+1)
for k, v := range hostVars {
scopedVars[k] = v
}
if item != nil {
scopedVars["item"] = item
}
out := make(map[string]interface{}, len(args))
for k, v := range args {
s, ok := v.(string)
if !ok || !strings.Contains(s, "{{") {
out[k] = v
continue
}
rendered, err := tmpl.Render("arg:"+k, s, scopedVars)
if err != nil {
return nil, err
}
out[k] = rendered
}
return out, nil
}