From 7b8647e8bc7786c7872da84ec7a72d13a0eee478 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Tue, 12 May 2026 19:19:46 -0700 Subject: [PATCH] fix(shell): loop s3.lifecycle.run-shard so CI workflow stays alive (#9476) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The s3tests workflow (.github/workflows/s3tests.yml) backgrounds `weed shell -c 's3.lifecycle.run-shard -shards 0-15 -s3 ... -refresh 2s'` and then runs `kill -0 $pid` to confirm the worker stayed alive. The PR-9475 restore ran dailyrun.Run once and exited cleanly — even faster when no buckets had lifecycle rules yet ("nothing to run"). The aliveness check then failed and the s3tests job died with "lifecycle worker died on startup". Caught on https://github.com/seaweedfs/seaweedfs/actions/runs/25772523143/job/75698413401. Fix: - -refresh now drives an inter-pass loop. cadence=0 (default) is one-shot, matching the test/s3/lifecycle/ integration-test invocation that omits -refresh and expects synchronous return. cadence>0 (the CI case) keeps the command alive until -runtime expires, running a fresh dailyrun.Run on every tick. - Each iteration re-loads bucket configs via scheduler.LoadCompileInputs so rules created mid-run (the s3tests flow creates rules AFTER the worker starts) get picked up. - The "no rules; nothing to run" early return is gone — the command stays alive even with an empty initial snapshot, waiting for tests to add rules. - -dispatch, -checkpoint, -bootstrap-interval stay accepted-but- ignored (legacy streaming flags). --- weed/shell/command_s3_lifecycle_run_shard.go | 134 +++++++++++-------- 1 file changed, 78 insertions(+), 56 deletions(-) diff --git a/weed/shell/command_s3_lifecycle_run_shard.go b/weed/shell/command_s3_lifecycle_run_shard.go index b11132641..51c383a18 100644 --- a/weed/shell/command_s3_lifecycle_run_shard.go +++ b/weed/shell/command_s3_lifecycle_run_shard.go @@ -2,6 +2,7 @@ package shell import ( "context" + "errors" "flag" "fmt" "io" @@ -65,15 +66,18 @@ func (c *commandS3LifecycleRunShard) Do(args []string, env *CommandEnv, writer i shard := fs.Int("shard", -1, "single shard id in [0, 16); use -shards for a range or set") shardsSpec := fs.String("shards", "", "shard range \"lo-hi\" or comma list \"a,b,c\"; mutually exclusive with -shard") s3Endpoint := fs.String("s3", "", "s3 server gRPC endpoint, host:port") - eventBudget := fs.Int("events", 1000, "max in-shard events to process before returning (0 = unbounded)") - runtime := fs.Duration("runtime", 0, "wall-clock cap on the run; 0 = no timeout") + eventBudget := fs.Int("events", 1000, "max in-shard events per pass (0 = drain to now)") + runtime := fs.Duration("runtime", 0, "wall-clock cap on the whole run; 0 = no timeout") + // -refresh drives the inter-pass cadence when the command runs as a + // long-lived worker (the s3tests CI workflow case): every refresh + // the engine snapshot is re-loaded and dailyrun.Run fires another + // pass. 0 means "run once and exit" (the integration-test case). + cadence := fs.Duration("refresh", 0, "inter-pass interval; 0 = single pass, then exit") // Obsolete flags kept for back-compat with existing CI scripts and - // integration tests. The daily-replay path has no streaming-tick - // cadence; accept and ignore. - _ = fs.Duration("dispatch", 0, "ignored (legacy streaming-dispatcher flag)") - _ = fs.Duration("checkpoint", 0, "ignored (legacy streaming-checkpoint flag)") - _ = fs.Duration("refresh", 0, "ignored (legacy streaming-refresh flag)") - _ = fs.Duration("bootstrap-interval", 0, "ignored (legacy streaming-bootstrap flag)") + // integration tests. Accept and ignore. + _ = fs.Duration("dispatch", 0, "ignored (legacy streaming flag)") + _ = fs.Duration("checkpoint", 0, "ignored (legacy streaming flag)") + _ = fs.Duration("bootstrap-interval", 0, "ignored (legacy streaming flag)") if err := fs.Parse(args); err != nil { return err } @@ -105,61 +109,79 @@ func (c *commandS3LifecycleRunShard) Do(args []string, env *CommandEnv, writer i rpcClient := s3_lifecycle_pb.NewSeaweedS3LifecycleInternalClient(conn) return env.WithFilerClient(true, func(filerClient filer_pb.SeaweedFilerClient) error { - eng := engine.New() - inputs, parseErrors, err := scheduler.LoadCompileInputs(context.Background(), filerClient, bucketsPath) - if err != nil { - return fmt.Errorf("load lifecycle configs: %w", err) - } - for i, pe := range parseErrors { - if i < 3 { - fmt.Fprintf(writer, "warning: %s: %v\n", pe.Bucket, pe.Err) - } - } - if extra := len(parseErrors) - 3; extra > 0 { - fmt.Fprintf(writer, "warning: %d additional bucket(s) had malformed lifecycle config\n", extra) - } - eng.Compile(inputs, engine.CompileOptions{PriorStates: scheduler.AllActivePriorStates(inputs)}) - if len(inputs) == 0 { - fmt.Fprintln(writer, "no buckets with enabled lifecycle rules; nothing to run") - return nil - } - - buckets := make([]string, 0, len(inputs)) - for _, in := range inputs { - if in.Bucket != "" { - buckets = append(buckets, in.Bucket) - } - } - client := &lifecycleClientCallable{c: rpcClient} - listFn := dailyrun.FilerListFunc(filerClient, bucketsPath) - walkerDispatch := &dailyrun.WalkerDispatcher{Client: client} - walker := dailyrun.WalkerFunc(func(walkCtx context.Context, view *engine.Snapshot, shardID int) error { - return dailyrun.WalkBuckets(walkCtx, view, shardID, buckets, listFn, walkerDispatch) - }) - ctx := context.Background() var cancel context.CancelFunc if *runtime > 0 { ctx, cancel = context.WithTimeout(ctx, *runtime) defer cancel() } + client := &lifecycleClientCallable{c: rpcClient} + listFn := dailyrun.FilerListFunc(filerClient, bucketsPath) + walkerDispatch := &dailyrun.WalkerDispatcher{Client: client} - fmt.Fprintf(writer, "running shards %s (event budget=%d, runtime=%s)…\n", - formatShardLabel(shards), *eventBudget, *runtime) - runErr := dailyrun.Run(ctx, dailyrun.Config{ - Shards: shards, - BucketsPath: bucketsPath, - Engine: eng, - FilerClient: filerClient, - Client: client, - Persister: &dailyrun.FilerCursorPersister{Store: dispatcher.NewFilerStoreClient(filerClient)}, - Lister: dispatcher.NewFilerSiblingLister(filerClient, bucketsPath), - Walker: walker, - EventBudget: *eventBudget, - ClientName: fmt.Sprintf("shell-lifecycle-%s", formatShardLabel(shards)), - }) - if runErr != nil { - return fmt.Errorf("daily_run: %w", runErr) + fmt.Fprintf(writer, "running shards %s (event budget=%d, runtime=%s, refresh=%s)…\n", + formatShardLabel(shards), *eventBudget, *runtime, *cadence) + + runPass := func() error { + eng := engine.New() + inputs, parseErrors, err := scheduler.LoadCompileInputs(ctx, filerClient, bucketsPath) + if err != nil { + return fmt.Errorf("load lifecycle configs: %w", err) + } + for i, pe := range parseErrors { + if i < 3 { + fmt.Fprintf(writer, "warning: %s: %v\n", pe.Bucket, pe.Err) + } + } + if extra := len(parseErrors) - 3; extra > 0 { + fmt.Fprintf(writer, "warning: %d additional bucket(s) had malformed lifecycle config\n", extra) + } + eng.Compile(inputs, engine.CompileOptions{PriorStates: scheduler.AllActivePriorStates(inputs)}) + if len(inputs) == 0 { + return nil + } + buckets := make([]string, 0, len(inputs)) + for _, in := range inputs { + if in.Bucket != "" { + buckets = append(buckets, in.Bucket) + } + } + walker := dailyrun.WalkerFunc(func(walkCtx context.Context, view *engine.Snapshot, shardID int) error { + return dailyrun.WalkBuckets(walkCtx, view, shardID, buckets, listFn, walkerDispatch) + }) + return dailyrun.Run(ctx, dailyrun.Config{ + Shards: shards, + BucketsPath: bucketsPath, + Engine: eng, + FilerClient: filerClient, + Client: client, + Persister: &dailyrun.FilerCursorPersister{Store: dispatcher.NewFilerStoreClient(filerClient)}, + Lister: dispatcher.NewFilerSiblingLister(filerClient, bucketsPath), + Walker: walker, + EventBudget: *eventBudget, + ClientName: fmt.Sprintf("shell-lifecycle-%s", formatShardLabel(shards)), + }) + } + + if err := runPass(); err != nil && !errors.Is(err, context.Canceled) && !errors.Is(err, context.DeadlineExceeded) { + return fmt.Errorf("daily_run: %w", err) + } + // cadence=0 → one-shot (test/s3/lifecycle/ uses this). + // cadence>0 → loop until ctx done (s3tests CI workflow uses this). + if *cadence > 0 { + ticker := time.NewTicker(*cadence) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + fmt.Fprintf(writer, "shards %s complete; ctx done\n", formatShardLabel(shards)) + return nil + case <-ticker.C: + if err := runPass(); err != nil && !errors.Is(err, context.Canceled) && !errors.Is(err, context.DeadlineExceeded) { + fmt.Fprintf(writer, "shards %s pass error: %v\n", formatShardLabel(shards), err) + } + } + } } fmt.Fprintf(writer, "shards %s complete; cursors checkpointed\n", formatShardLabel(shards)) return nil