From 167a1e66b438e6b412974be55de6a8fede85a220 Mon Sep 17 00:00:00 2001 From: Ryan Cook Date: Fri, 9 Oct 2026 11:17:31 -0400 Subject: [PATCH] Fix graceful engine shutdown on SIGTERM and SIGINT Signed-off-by: Ryan Cook --- .github/workflows/docker-image.yml | 3 + .github/workflows/unit-tests.yml | 2 +- docs/release_notes.rst | 5 ++ docs/running.rst | 12 ++++ pkg/engine/fetchit.go | 48 ++++++++++--- pkg/engine/shutdown.go | 77 ++++++++++++++++++++ pkg/engine/shutdown_test.go | 111 +++++++++++++++++++++++++++++ pkg/engine/start.go | 12 ++-- pkg/engine/types.go | 4 +- scripts/entry.sh | 4 +- scripts/test-shutdown.sh | 60 ++++++++++++++++ 11 files changed, 319 insertions(+), 19 deletions(-) create mode 100644 pkg/engine/shutdown.go create mode 100644 pkg/engine/shutdown_test.go create mode 100755 scripts/test-shutdown.sh diff --git a/.github/workflows/docker-image.yml b/.github/workflows/docker-image.yml index 530f70a9..87a049ea 100644 --- a/.github/workflows/docker-image.yml +++ b/.github/workflows/docker-image.yml @@ -212,6 +212,9 @@ jobs: - name: tag the image run: sudo podman tag quay.io/fetchit/fetchit-amd:latest quay.io/fetchit/fetchit:latest + - name: Verify graceful engine shutdown + run: sudo bash scripts/test-shutdown.sh quay.io/fetchit/fetchit-amd:latest + - name: create fetchit config directory run: sudo mkdir /root/.fetchit diff --git a/.github/workflows/unit-tests.yml b/.github/workflows/unit-tests.yml index e45dbfb1..6c413e5b 100644 --- a/.github/workflows/unit-tests.yml +++ b/.github/workflows/unit-tests.yml @@ -36,7 +36,7 @@ jobs: AGE_TEST_BINARY: /usr/bin/age-keygen run: go test -mod=readonly -tags "$BUILDTAGS" ./... - name: Engine concurrency race checks - run: go test -race -mod=readonly -tags "$BUILDTAGS" ./pkg/engine -run 'TestStatus|TestMirror|TestCommonGit|TestHelperContainer|TestRetirement|TestConfigReloadCan|TestRollback|TestDisconnected|TestFileTransfer' + run: go test -race -mod=readonly -tags "$BUILDTAGS" ./pkg/engine -run 'TestStatus|TestMirror|TestCommonGit|TestHelperContainer|TestRetirement|TestConfigReloadCan|TestRollback|TestDisconnected|TestFileTransfer|TestShutdown' fedora: runs-on: ubuntu-26.04 diff --git a/docs/release_notes.rst b/docs/release_notes.rst index eb9bfffa..51feb272 100644 --- a/docs/release_notes.rst +++ b/docs/release_notes.rst @@ -4,6 +4,11 @@ Release notes Unreleased ---------- +* Graceful engine shutdown (issue #287): SIGTERM and SIGINT stop scheduling + work and cancel running operations, with a five-second grace period before + exit. The entry script forwards signals by replacing itself with FetchIt. + Deployed workloads remain running. See :doc:`running`. + * Per-declaration named-volume creation: raw mounts support ``create: true``; kube workloads can generate missing PVC declarations with the ``fetchit.containers.io/create-volumes`` annotation. Existing volumes and data diff --git a/docs/running.rst b/docs/running.rst index cc00cd65..71c66556 100644 --- a/docs/running.rst +++ b/docs/running.rst @@ -7,6 +7,18 @@ rootless Podman stores are separate: build/pull helper images and create network in the same store used by FetchIt. See :doc:`samples` for amd64/arm64 applications and :doc:`release_notes` for features requiring a newer engine/helper image. +Stopping the engine +------------------- + +On current main, ``podman stop fetchit`` sends SIGTERM directly to the engine. +SIGINT is also supported. FetchIt stops scheduling work, cancels running +operations, and waits up to five seconds for jobs to finish before exiting. +Operations that do not honor cancellation may still be in flight when that +grace period expires; check the logs and workload state before restarting. +Stopping FetchIt leaves its deployed workloads running and preserves engine +state for the next start. Published releases require this shutdown fix to be +included in their engine image. + Rootless socket --------------- diff --git a/pkg/engine/fetchit.go b/pkg/engine/fetchit.go index a5c9a917..8605070a 100644 --- a/pkg/engine/fetchit.go +++ b/pkg/engine/fetchit.go @@ -41,6 +41,7 @@ type Fetchit struct { runMu sync.RWMutex retired bool runCancel context.CancelFunc + lifetime context.Context removals *removalStore // conn holds podman client conn context.Context @@ -91,11 +92,20 @@ func Execute() { func (fc *FetchitConfig) Restart() { configRestartMu.Lock() defer configRestartMu.Unlock() + if fc.runtimeContext().Err() != nil { + return + } old := fetchit old.retire() old.scheduler.Clear() + if fc.runtimeContext().Err() != nil { + return + } next := fc.InitConfig(false) - cobra.CheckErr(next.startTargets()) + err := next.startTargets() + if fc.runtimeContext().Err() == nil { + cobra.CheckErr(err) + } } // Config reload runs outside this gate so it can wait for deployment jobs to @@ -110,7 +120,7 @@ func (f *Fetchit) retire() { } func (f *Fetchit) runMethod(m Method, ctx, conn context.Context, skew int) { f.runMu.RLock() - if f.retired { + if f.retired || ctx.Err() != nil { f.runMu.RUnlock() return } @@ -123,7 +133,7 @@ func (f *Fetchit) runMethod(m Method, ctx, conn context.Context, skew int) { f.runMu.RLock() retired := f.retired f.runMu.RUnlock() - if retired { + if retired || ctx.Err() != nil { return } status.recordRun(m) @@ -157,20 +167,28 @@ func readConfig(v *viper.Viper) (*FetchitConfig, bool, error) { func (fc *FetchitConfig) populateFetchit(config *FetchitConfig) *Fetchit { fetchit = newFetchit() - ctx := context.Background() + ctx := fc.runtimeContext() if fc.conn == nil { // TODO: socket directory same for all platforms? // sock_dir := os.Getenv("XDG_RUNTIME_DIR") // socket := "unix:" + sock_dir + "/podman/podman.sock" conn, err := bindings.NewConnection(ctx, "unix://run/podman/podman.sock") + if ctx.Err() != nil { + fetchit.lifetime = ctx + return fetchit + } if err != nil || conn == nil { cobra.CheckErr(fmt.Errorf("error establishing connection to podman.sock: %v", err)) } fc.conn = conn } fetchit.conn = fc.conn + fetchit.lifetime = ctx if err := detectOrFetchImage(fc.conn, fetchitImage, false); err != nil { + if ctx.Err() != nil { + return fetchit + } cobra.CheckErr(err) } @@ -400,12 +418,14 @@ func getMethodTargetScheds(targetConfigs []*TargetConfig, fetchit *Fetchit) *Fet return fetchit } -func (f *Fetchit) RunTargets() { - cobra.CheckErr(f.startTargets()) - select {} -} - func (f *Fetchit) startTargets() error { + parent := f.lifetime + if parent == nil { + parent = context.Background() + } + if err := parent.Err(); err != nil { + return err + } store, err := loadRemovalStore(defaultRemovalsPath) if err != nil { return err @@ -417,11 +437,14 @@ func (f *Fetchit) startTargets() error { if err := store.reconcile(f.conn); err != nil { logger.Errorf("Method removal cleanup: %v", err) } - ctx, cancel := context.WithCancel(context.Background()) + ctx, cancel := context.WithCancel(parent) f.runCancel = cancel status.replace(f.methodTargetScheds) startStatusServer() for method := range f.methodTargetScheds { + if err := ctx.Err(); err != nil { + return err + } // ConfigReload, PodmanAutoUpdateAll, Image, Prune methods do not include git URL if method.GetTarget().url != "" { if err := getRepo(method.GetTarget()); err != nil { @@ -429,6 +452,9 @@ func (f *Fetchit) startTargets() error { } } } + if err := ctx.Err(); err != nil { + return err + } s := f.scheduler for method, schedInfo := range f.methodTargetScheds { @@ -447,7 +473,7 @@ func (f *Fetchit) startTargets() error { if _, err := s.Every(1).Minute().Tag("removal-cleanup").Do(func() { f.runMu.RLock() defer f.runMu.RUnlock() - if f.retired { + if f.retired || ctx.Err() != nil { return } if err := store.reconcile(f.conn); err != nil { diff --git a/pkg/engine/shutdown.go b/pkg/engine/shutdown.go new file mode 100644 index 00000000..95d36bc4 --- /dev/null +++ b/pkg/engine/shutdown.go @@ -0,0 +1,77 @@ +package engine + +import ( + "context" + "fmt" + "os" + "time" +) + +// Leave room for Podman's default ten-second stop timeout. +const shutdownGracePeriod = 5 * time.Second + +func (fc *FetchitConfig) runtimeContext() context.Context { + if fc.lifetime != nil { + return fc.lifetime + } + return context.Background() +} + +func (fc *FetchitConfig) run(ctx context.Context) error { + ctx, cancel := context.WithCancel(ctx) + defer cancel() + fc.lifetime = ctx + started := make(chan error, 1) + go func() { + // Shutdown and reload must not race initial scheduler creation/startup. + configRestartMu.Lock() + defer configRestartMu.Unlock() + if err := ctx.Err(); err != nil { + started <- err + return + } + f := fc.InitConfig(true) + started <- f.startTargets() + }() + + select { + case err := <-started: + if err != nil && ctx.Err() == nil { + cancel() + fc.shutdown(shutdownGracePeriod) + return err + } + if err == nil { + <-ctx.Done() + } + case <-ctx.Done(): + } + fc.shutdown(shutdownGracePeriod) + return nil +} + +// The lifetime context is canceled before this call, so queued methods and +// reloads cannot start new work. Stop waits for scheduled jobs to finish, but +// initialization, Git operations, or helpers may not honor cancellation yet. +func (fc *FetchitConfig) shutdown(grace time.Duration) { + done := make(chan struct{}) + go func() { + configRestartMu.Lock() + scheduler := fc.scheduler + configRestartMu.Unlock() + // Never hold the reload mutex while waiting for jobs: a reload job may + // itself be waiting for that mutex before observing cancellation. + if scheduler != nil { + scheduler.Stop() + } + close(done) + }() + timer := time.NewTimer(grace) + defer timer.Stop() + select { + case <-done: + fmt.Fprintln(os.Stderr, "FetchIt shutdown complete") + case <-timer.C: + fmt.Fprintln(os.Stderr, "FetchIt shutdown grace period expired; exiting with work still in flight") + } +} diff --git a/pkg/engine/shutdown_test.go b/pkg/engine/shutdown_test.go new file mode 100644 index 00000000..57a3525c --- /dev/null +++ b/pkg/engine/shutdown_test.go @@ -0,0 +1,111 @@ +package engine + +import ( + "context" + "testing" + "time" + + "github.com/go-co-op/gocron" +) + +func waitShutdownTest(t *testing.T, ch <-chan struct{}) { + t.Helper() + select { + case <-ch: + case <-time.After(2 * time.Second): + t.Fatal("timed out waiting for shutdown test job") + } +} + +func TestShutdownCancelsAndWaitsForRunningJob(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + s := gocron.NewScheduler(time.UTC) + fc := &FetchitConfig{lifetime: ctx, scheduler: s} + f := newFetchit() + entered, canceled, release := make(chan struct{}), make(chan struct{}), make(chan struct{}) + defer close(release) + m := &lifecycleFake{CommonMethod: CommonMethod{target: &Target{}}, kind: rawMethod, process: func() { + close(entered) + <-ctx.Done() + close(canceled) + <-release + }} + if _, err := s.Every(1).Hour().Do(func() { f.runMethod(m, ctx, context.Background(), 0) }); err != nil { + t.Fatal(err) + } + s.StartAsync() + waitShutdownTest(t, entered) + cancel() + done := make(chan struct{}) + go func() { fc.shutdown(time.Second); close(done) }() + waitShutdownTest(t, canceled) + select { + case <-done: + t.Fatal("shutdown returned before running job finished") + case <-time.After(20 * time.Millisecond): + } + release <- struct{}{} + waitShutdownTest(t, done) + if s.IsRunning() { + t.Fatal("scheduler still running after shutdown") + } +} + +func TestShutdownBoundsUncooperativeJob(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + s := gocron.NewScheduler(time.UTC) + fc := &FetchitConfig{lifetime: ctx, scheduler: s} + entered, release, finished := make(chan struct{}), make(chan struct{}), make(chan struct{}) + defer close(release) + if _, err := s.Every(1).Hour().Do(func() { close(entered); <-release; close(finished) }); err != nil { + t.Fatal(err) + } + s.StartAsync() + waitShutdownTest(t, entered) + cancel() + started := time.Now() + fc.shutdown(40 * time.Millisecond) + if elapsed := time.Since(started); elapsed < 35*time.Millisecond || elapsed > time.Second { + t.Fatalf("shutdown did not respect grace period: %s", elapsed) + } + release <- struct{}{} + waitShutdownTest(t, finished) +} + +func TestShutdownRejectsQueuedMethodsAndReloads(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + f := newFetchit() + for _, kind := range []string{rawMethod, configFileMethod} { + m := &lifecycleFake{CommonMethod: CommonMethod{target: &Target{}}, kind: kind, process: func() { + t.Fatal("canceled job ran") + }} + f.runMethod(m, ctx, context.Background(), 0) + } + // A canceled reload must return before accessing/replacing the active engine. + fc := &FetchitConfig{lifetime: ctx} + fc.Restart() +} + +func TestShutdownBoundsStartupAndReloadLock(t *testing.T) { + configRestartMu.Lock() + defer configRestartMu.Unlock() + fc := &FetchitConfig{} + done := make(chan struct{}) + go func() { fc.shutdown(40 * time.Millisecond); close(done) }() + waitShutdownTest(t, done) +} + +func TestShutdownBeforeInitialization(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + fc := newFetchitConfig() + if err := fc.run(ctx); err != nil { + t.Fatalf("canceled startup returned an error: %v", err) + } + if fc.scheduler != nil || fc.conn != nil { + t.Fatal("canceled startup initialized the engine") + } +} diff --git a/pkg/engine/start.go b/pkg/engine/start.go index 12673ae1..47c09d95 100644 --- a/pkg/engine/start.go +++ b/pkg/engine/start.go @@ -1,11 +1,14 @@ package engine import ( + "os" + "os/signal" + "syscall" + "github.com/natefinch/lumberjack" "github.com/spf13/cobra" "go.uber.org/zap" "go.uber.org/zap/zapcore" - "os" ) // This file will be created within the fetchit pod @@ -15,9 +18,10 @@ var startCmd = &cobra.Command{ Use: "start", Short: "Start fetchit engine", Long: `Start fetchit engine`, - Run: func(cmd *cobra.Command, args []string) { - fetchit = fetchitConfig.InitConfig(true) - fetchit.RunTargets() + RunE: func(cmd *cobra.Command, args []string) error { + ctx, stop := signal.NotifyContext(cmd.Context(), syscall.SIGTERM, os.Interrupt) + defer stop() + return fetchitConfig.run(ctx) }, } diff --git a/pkg/engine/types.go b/pkg/engine/types.go index c055f8f0..06bb28c8 100644 --- a/pkg/engine/types.go +++ b/pkg/engine/types.go @@ -27,7 +27,9 @@ type FetchitConfig struct { PodmanAutoUpdate *PodmanAutoUpdate `mapstructure:"podmanAutoUpdate"` Images []*Image `mapstructure:"images"` conn context.Context - scheduler *gocron.Scheduler + // lifetime survives configuration reloads and is canceled on engine shutdown. + lifetime context.Context + scheduler *gocron.Scheduler } type TargetConfig struct { diff --git a/scripts/entry.sh b/scripts/entry.sh index 22069acf..e032eb11 100755 --- a/scripts/entry.sh +++ b/scripts/entry.sh @@ -20,7 +20,7 @@ main() { fi echo 'starting fetchit' - /usr/local/bin/fetchit start + exec /usr/local/bin/fetchit start } -main \ No newline at end of file +main diff --git a/scripts/test-shutdown.sh b/scripts/test-shutdown.sh new file mode 100755 index 00000000..498978a1 --- /dev/null +++ b/scripts/test-shutdown.sh @@ -0,0 +1,60 @@ +#!/bin/bash +# Run against a disposable Linux Podman host with its API socket enabled. +set -euo pipefail + +image=${1:-quay.io/fetchit/fetchit:latest} +socket=${PODMAN_SOCKET:-/run/podman/podman.sock} +name="fetchit-shutdown-$$" +config=$(mktemp -d) +cleanup() { + podman rm -f "$name" >/dev/null 2>&1 || true + rm -rf "$config" +} +trap cleanup EXIT +printf 'targetConfigs: []\n' > "$config/config.yaml" + +for signal in TERM INT; do + podman run -d --name "$name" \ + -v "$config:/opt/mount" \ + -v "$socket:/run/podman/podman.sock" \ + --security-opt label=disable \ + -e FETCHIT_STATUS_ADDR=:8080 \ + -e 'FETCHIT_CONFIG=targetConfigs: []' \ + -p 127.0.0.1::8080 "$image" + port=$(podman port "$name" 8080/tcp) + ready=false + for attempt in {1..60}; do + if curl --fail --silent --max-time 1 "http://$port/healthz" >/dev/null; then + ready=true + break + fi + sleep 0.5 + done + if [[ "$ready" != true ]]; then + podman logs "$name" + echo 'engine did not become ready' >&2 + exit 1 + fi + # A Bash wrapper at PID 1 was the cause of issue #287. + process=$(podman top "$name" pid,args) + if [[ "$process" == *entry.sh* ]]; then + echo "entry script still wraps FetchIt: $process" >&2 + exit 1 + fi + started=$SECONDS + if [[ "$signal" == TERM ]]; then + podman container stop "$name" + else + podman kill --signal INT "$name" + timeout 8 podman wait "$name" + fi + elapsed=$((SECONDS - started)) + exit_code=$(podman inspect --format '{{.State.ExitCode}}' "$name") + podman logs "$name" + if (( elapsed >= 8 )) || [[ "$exit_code" != 0 ]]; then + echo "$signal shutdown failed: ${elapsed}s, exit $exit_code" >&2 + exit 1 + fi + echo "$signal shutdown passed: ${elapsed}s, exit $exit_code" + podman rm "$name" +done