Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .github/workflows/docker-image.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/unit-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions docs/release_notes.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 12 additions & 0 deletions docs/running.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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
---------------

Expand Down
48 changes: 37 additions & 11 deletions pkg/engine/fetchit.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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
}
Expand All @@ -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)
Expand Down Expand Up @@ -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)
}

Expand Down Expand Up @@ -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
Expand All @@ -417,18 +437,24 @@ 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 {
logger.Debugf("Target: %s, clone error: %v, will retry next scheduled run", method.GetTarget(), err)
}
}
}
if err := ctx.Err(); err != nil {
return err
}

s := f.scheduler
for method, schedInfo := range f.methodTargetScheds {
Expand All @@ -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 {
Expand Down
77 changes: 77 additions & 0 deletions pkg/engine/shutdown.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
111 changes: 111 additions & 0 deletions pkg/engine/shutdown_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
Loading
Loading