|
| 1 | +// Package hooks runs the declarative lifecycle hooks (spec 11) at saga phase |
| 2 | +// boundaries: user commands on the host (run: host) or inside a running service |
| 3 | +// (run: exec via `compose exec -T`), with per-hook timeout, retries, and |
| 4 | +// onFailure semantics. firstRun/postPull hooks are made idempotent by the |
| 5 | +// `hook_run` ledger — recorded ONLY on success and ONLY inside the flock, so a |
| 6 | +// failed or interrupted hook re-runs next time (the correct replacement for |
| 7 | +// Postgres initdb.d, DECISIONS D8). |
| 8 | +// |
| 9 | +// The hook BODY runs OUTSIDE the global flock (a 10-minute `npm install` must not |
| 10 | +// serialize every other invocation, spec 08 #1 rule); only the ledger record is |
| 11 | +// taken under the lock. This is the thin runner — ${ref}/${env}/${self}/secret:// |
| 12 | +// interpolation happens upstream in the saga (the resolved values arrive via |
| 13 | +// PhaseOpts.ExtraEnv); the full workspace-scope ordering is X3. |
| 14 | +package hooks |
| 15 | + |
| 16 | +import ( |
| 17 | + "context" |
| 18 | + "fmt" |
| 19 | + "os" |
| 20 | + "os/exec" |
| 21 | + "path/filepath" |
| 22 | + "sort" |
| 23 | + "time" |
| 24 | + |
| 25 | + "github.com/open-source-cloud/devstack/internal/config" |
| 26 | +) |
| 27 | + |
| 28 | +// Defaults (spec 11). |
| 29 | +const ( |
| 30 | + DefaultTimeout = 120 * time.Second |
| 31 | + defaultBackoff = 2 * time.Second |
| 32 | +) |
| 33 | + |
| 34 | +// onFailure modes (spec 11). |
| 35 | +const ( |
| 36 | + OnAbort = "abort" |
| 37 | + OnWarn = "warn" |
| 38 | + OnContinue = "continue" |
| 39 | +) |
| 40 | + |
| 41 | +// Result statuses for a single hook. |
| 42 | +const ( |
| 43 | + StatusRan = "ran" // executed successfully |
| 44 | + StatusSkipped = "skipped" // ledger already satisfied (idempotent hook) |
| 45 | + StatusWarned = "warned" // failed but onFailure was warn/continue |
| 46 | + StatusFailed = "failed" // failed with onFailure=abort |
| 47 | +) |
| 48 | + |
| 49 | +// Execer runs the two hook transports, returning combined stdout+stderr. It is |
| 50 | +// injectable so the runner is testable without a daemon. Implementations MUST |
| 51 | +// honor ctx cancellation (the runner enforces each hook's timeout via ctx). |
| 52 | +type Execer interface { |
| 53 | + // Host runs argv on the host. workdir is hook-relative (joined under the |
| 54 | + // documented base dir); env is KEY=VALUE pairs appended to the process env. |
| 55 | + Host(ctx context.Context, workdir string, env, argv []string) (string, error) |
| 56 | + // Exec runs argv inside a RUNNING service via `compose exec -T`. workdir is the |
| 57 | + // in-container -w path; env becomes repeated -e NAME=VALUE flags. |
| 58 | + Exec(ctx context.Context, service, workdir string, env, argv []string) (string, error) |
| 59 | +} |
| 60 | + |
| 61 | +// OSExecer is the production Execer: os/exec for host hooks, the `docker compose |
| 62 | +// exec` CLI for service hooks (DECISIONS D5 — devstack owns compose CLI |
| 63 | +// construction). -T disables TTY for deterministic non-interactive runs; combined |
| 64 | +// output is captured for the saga checklist and the failure remediation. |
| 65 | +type OSExecer struct { |
| 66 | + BaseDir string // documented working dir for run:host (repo root / workspace root) |
| 67 | + Project string // compose -p |
| 68 | + File string // compose -f |
| 69 | +} |
| 70 | + |
| 71 | +// Host runs argv on the host from the hook's working directory. |
| 72 | +func (e OSExecer) Host(ctx context.Context, workdir string, env, argv []string) (string, error) { |
| 73 | + if len(argv) == 0 { |
| 74 | + return "", fmt.Errorf("empty command") |
| 75 | + } |
| 76 | + dir := e.BaseDir |
| 77 | + if workdir != "" { |
| 78 | + if filepath.IsAbs(workdir) { |
| 79 | + dir = workdir |
| 80 | + } else { |
| 81 | + dir = filepath.Join(e.BaseDir, workdir) |
| 82 | + } |
| 83 | + } |
| 84 | + cmd := exec.CommandContext(ctx, argv[0], argv[1:]...) |
| 85 | + cmd.Dir = dir |
| 86 | + cmd.Env = append(os.Environ(), env...) |
| 87 | + out, err := cmd.CombinedOutput() |
| 88 | + return string(out), err |
| 89 | +} |
| 90 | + |
| 91 | +// Exec shells argv into a running service. compose exec -e accepts NAME=VALUE |
| 92 | +// (unlike `up`), so hook secrets are passed inline without ever touching a file. |
| 93 | +func (e OSExecer) Exec(ctx context.Context, service, workdir string, env, argv []string) (string, error) { |
| 94 | + if len(argv) == 0 { |
| 95 | + return "", fmt.Errorf("empty command") |
| 96 | + } |
| 97 | + args := []string{"compose", "-p", e.Project, "-f", e.File, "exec", "-T"} |
| 98 | + if workdir != "" { |
| 99 | + args = append(args, "-w", workdir) |
| 100 | + } |
| 101 | + for _, kv := range env { |
| 102 | + args = append(args, "-e", kv) |
| 103 | + } |
| 104 | + args = append(args, service) |
| 105 | + args = append(args, argv...) |
| 106 | + cmd := exec.CommandContext(ctx, "docker", args...) |
| 107 | + cmd.Dir = e.BaseDir |
| 108 | + out, err := cmd.CombinedOutput() |
| 109 | + return string(out), err |
| 110 | +} |
| 111 | + |
| 112 | +// Ledger is the idempotency store the runner needs (satisfied by *state.DB). Kept |
| 113 | +// as an interface so internal/hooks does not import internal/state. |
| 114 | +type Ledger interface { |
| 115 | + HookSatisfied(project, hook, scopeKey string) (bool, error) |
| 116 | + RecordHookRun(project, hook, scopeKey string) error |
| 117 | +} |
| 118 | + |
| 119 | +// Locker runs fn while holding the machine-global flock (wraps lock.WithLock). |
| 120 | +type Locker func(ctx context.Context, fn func() error) error |
| 121 | + |
| 122 | +// Runner executes lifecycle hooks. Execer is required; Ledger+Lock are needed |
| 123 | +// only for idempotent phases; Backoff/Logf are optional. |
| 124 | +type Runner struct { |
| 125 | + Execer Execer |
| 126 | + Ledger Ledger |
| 127 | + Lock Locker |
| 128 | + Backoff time.Duration // between retry attempts (0 → default 2s) |
| 129 | + Logf func(format string, a ...any) // optional progress/warn sink |
| 130 | +} |
| 131 | + |
| 132 | +// Result is the outcome of one hook (for the saga checklist / --json). |
| 133 | +type Result struct { |
| 134 | + Hook string |
| 135 | + Status string |
| 136 | + Attempts int |
| 137 | + Output string |
| 138 | + Err error |
| 139 | +} |
| 140 | + |
| 141 | +// Run executes one hook with its timeout and retries, returning combined output, |
| 142 | +// the attempt count, and the failure error (after retries) if any. It applies NO |
| 143 | +// onFailure or idempotency semantics — RunPhase layers those on. extraEnv carries |
| 144 | +// resolved ${ref}/secret values (saga-supplied), appended after the hook's own |
| 145 | +// env so secrets land last. |
| 146 | +func (r *Runner) Run(ctx context.Context, h config.Hook, extraEnv []string) (string, int, error) { |
| 147 | + if err := validate(h); err != nil { |
| 148 | + return "", 0, err |
| 149 | + } |
| 150 | + timeout := DefaultTimeout |
| 151 | + if d, err := time.ParseDuration(h.Timeout); err == nil && h.Timeout != "" { |
| 152 | + timeout = d |
| 153 | + } |
| 154 | + env := buildEnv(h.Env, extraEnv) |
| 155 | + |
| 156 | + var ( |
| 157 | + out string |
| 158 | + err error |
| 159 | + attempts int |
| 160 | + ) |
| 161 | + for attempt := 0; attempt <= max(h.Retries, 0); attempt++ { |
| 162 | + if attempt > 0 { |
| 163 | + if e := sleep(ctx, r.backoff()); e != nil { |
| 164 | + return out, attempts, e |
| 165 | + } |
| 166 | + } |
| 167 | + attempts++ |
| 168 | + actx, cancel := context.WithTimeout(ctx, timeout) |
| 169 | + if h.Run == "exec" { |
| 170 | + out, err = r.Execer.Exec(actx, h.Service, h.Workdir, env, h.Command) |
| 171 | + } else { |
| 172 | + out, err = r.Execer.Host(actx, h.Workdir, env, h.Command) |
| 173 | + } |
| 174 | + cancel() |
| 175 | + if err == nil { |
| 176 | + return out, attempts, nil |
| 177 | + } |
| 178 | + } |
| 179 | + return out, attempts, fmt.Errorf("hook %q failed after %d attempt(s): %w", h.Name, attempts, err) |
| 180 | +} |
| 181 | + |
| 182 | +// PhaseOpts configures a phase run. |
| 183 | +type PhaseOpts struct { |
| 184 | + Project string // ledger 'project' |
| 185 | + Phase string // ledger 'hook' column, e.g. "firstRun" |
| 186 | + Idempotent bool // guard every hook via the ledger (firstRun/postPull) |
| 187 | + ScopeKey func(config.Hook) string // scope_key for idempotent / `once` hooks |
| 188 | + DefaultOnFailure string // applied when a hook omits onFailure |
| 189 | + ExtraEnv []string // resolved env/secret values for every hook |
| 190 | +} |
| 191 | + |
| 192 | +// RunPhase runs the phase's hooks in order, applying idempotency and onFailure: |
| 193 | +// - abort : stop immediately, return the error (phase fails). |
| 194 | +// - warn : log and continue; the phase still succeeds. |
| 195 | +// - continue: run remaining hooks, but the phase ultimately fails. |
| 196 | +// |
| 197 | +// A successful run of an idempotent (or `once`) hook is recorded in the ledger |
| 198 | +// INSIDE the flock; nothing is recorded on failure (spec 11). |
| 199 | +func (r *Runner) RunPhase(ctx context.Context, hookList []config.Hook, o PhaseOpts) ([]Result, error) { |
| 200 | + var ( |
| 201 | + results []Result |
| 202 | + deferredErr error // from onFailure: continue |
| 203 | + ) |
| 204 | + for _, h := range hookList { |
| 205 | + guarded := o.Idempotent || h.Once |
| 206 | + var scope string |
| 207 | + if guarded { |
| 208 | + if o.ScopeKey == nil { |
| 209 | + return results, fmt.Errorf("hook %q: idempotent phase %q has no scope_key function", h.Name, o.Phase) |
| 210 | + } |
| 211 | + scope = o.ScopeKey(h) |
| 212 | + if r.Ledger != nil { |
| 213 | + ok, err := r.Ledger.HookSatisfied(o.Project, o.Phase, scope) |
| 214 | + if err != nil { |
| 215 | + return results, err |
| 216 | + } |
| 217 | + if ok { |
| 218 | + results = append(results, Result{Hook: h.Name, Status: StatusSkipped}) |
| 219 | + continue |
| 220 | + } |
| 221 | + } |
| 222 | + } |
| 223 | + |
| 224 | + out, attempts, err := r.Run(ctx, h, o.ExtraEnv) |
| 225 | + res := Result{Hook: h.Name, Attempts: attempts, Output: out, Err: err} |
| 226 | + if err == nil { |
| 227 | + res.Status = StatusRan |
| 228 | + results = append(results, res) |
| 229 | + if guarded { |
| 230 | + if e := r.record(ctx, o.Project, o.Phase, scope); e != nil { |
| 231 | + return results, e |
| 232 | + } |
| 233 | + } |
| 234 | + continue |
| 235 | + } |
| 236 | + |
| 237 | + switch onFailure(h, o.DefaultOnFailure) { |
| 238 | + case OnWarn: |
| 239 | + res.Status = StatusWarned |
| 240 | + r.warn("hook %q failed (onFailure=warn): %v", h.Name, err) |
| 241 | + results = append(results, res) |
| 242 | + case OnContinue: |
| 243 | + res.Status = StatusWarned |
| 244 | + deferredErr = err |
| 245 | + r.warn("hook %q failed (onFailure=continue): %v", h.Name, err) |
| 246 | + results = append(results, res) |
| 247 | + default: // abort |
| 248 | + res.Status = StatusFailed |
| 249 | + results = append(results, res) |
| 250 | + return results, err |
| 251 | + } |
| 252 | + } |
| 253 | + return results, deferredErr |
| 254 | +} |
| 255 | + |
| 256 | +func (r *Runner) record(ctx context.Context, project, phase, scope string) error { |
| 257 | + if r.Ledger == nil { |
| 258 | + return nil |
| 259 | + } |
| 260 | + rec := func() error { return r.Ledger.RecordHookRun(project, phase, scope) } |
| 261 | + if r.Lock != nil { |
| 262 | + return r.Lock(ctx, rec) |
| 263 | + } |
| 264 | + return rec() |
| 265 | +} |
| 266 | + |
| 267 | +func (r *Runner) backoff() time.Duration { |
| 268 | + if r.Backoff > 0 { |
| 269 | + return r.Backoff |
| 270 | + } |
| 271 | + return defaultBackoff |
| 272 | +} |
| 273 | + |
| 274 | +func (r *Runner) warn(format string, a ...any) { |
| 275 | + if r.Logf != nil { |
| 276 | + r.Logf(format, a...) |
| 277 | + } |
| 278 | +} |
| 279 | + |
| 280 | +// validate enforces the transport-specific rule config can't (spec 11): run:exec |
| 281 | +// needs a target service. (Structural shape — name/run/command — is already |
| 282 | +// validated at config load, C3a.) |
| 283 | +func validate(h config.Hook) error { |
| 284 | + if len(h.Command) == 0 { |
| 285 | + return fmt.Errorf("hook %q: command is empty", h.Name) |
| 286 | + } |
| 287 | + if h.Run == "exec" && h.Service == "" { |
| 288 | + return fmt.Errorf("hook %q: run: exec requires a `service`", h.Name) |
| 289 | + } |
| 290 | + return nil |
| 291 | +} |
| 292 | + |
| 293 | +func onFailure(h config.Hook, def string) string { |
| 294 | + if h.OnFailure != "" { |
| 295 | + return h.OnFailure |
| 296 | + } |
| 297 | + if def != "" { |
| 298 | + return def |
| 299 | + } |
| 300 | + return OnAbort |
| 301 | +} |
| 302 | + |
| 303 | +// buildEnv renders the hook's env map (sorted for determinism) followed by the |
| 304 | +// saga-supplied extra env (secrets last, spec 11). |
| 305 | +func buildEnv(m map[string]string, extra []string) []string { |
| 306 | + keys := make([]string, 0, len(m)) |
| 307 | + for k := range m { |
| 308 | + keys = append(keys, k) |
| 309 | + } |
| 310 | + sort.Strings(keys) |
| 311 | + env := make([]string, 0, len(m)+len(extra)) |
| 312 | + for _, k := range keys { |
| 313 | + env = append(env, k+"="+m[k]) |
| 314 | + } |
| 315 | + return append(env, extra...) |
| 316 | +} |
| 317 | + |
| 318 | +// sleep is the cancelable inter-retry wait, indirected for tests. |
| 319 | +var sleep = func(ctx context.Context, d time.Duration) error { |
| 320 | + if d <= 0 { |
| 321 | + return ctx.Err() |
| 322 | + } |
| 323 | + t := time.NewTimer(d) |
| 324 | + defer t.Stop() |
| 325 | + select { |
| 326 | + case <-ctx.Done(): |
| 327 | + return ctx.Err() |
| 328 | + case <-t.C: |
| 329 | + return nil |
| 330 | + } |
| 331 | +} |
0 commit comments