From 2237a640d2400e337d9e6c638a7b9534cd9c77b9 Mon Sep 17 00:00:00 2001 From: Jakub Zych Date: Tue, 29 Sep 2026 22:34:21 +0200 Subject: [PATCH] feat(11-02): add schedule:run as a scheduler process and a cron --once mode - schedule:run runs a worker on the scheduled queue with every plugin's periodic jobs - schedule:run --once runs entries due in the current app.timezone minute without River, warns and skips unregistered commands, and returns the first command error - summer schedule:run delegate forwards --once - tests for --once minute matching, forged-entry skipping (T-11-09) and ByPeriod dedupe --- cmd/summer/main.go | 3 +- cmd/summer/main_test.go | 3 +- cmd/summer/runtime.go | 25 ++- modules/conga/README.md | 7 +- modules/conga/commands.go | 111 +++++++++- modules/conga/schedule_test.go | 357 +++++++++++++++++++++++++++++++++ modules/conga/scheduler.go | 33 ++- 7 files changed, 521 insertions(+), 18 deletions(-) diff --git a/cmd/summer/main.go b/cmd/summer/main.go index 2cc8b2c..ea1d546 100644 --- a/cmd/summer/main.go +++ b/cmd/summer/main.go @@ -6,9 +6,9 @@ import ( "os" "strings" - "git.golem15.com/golem15/summercms/modules/bonfire" "git.golem15.com/golem15/summercms/internal/build" "git.golem15.com/golem15/summercms/internal/dev" + "git.golem15.com/golem15/summercms/modules/bonfire" ) func main() { @@ -43,6 +43,7 @@ func toolCommands() []bonfire.Command { delegateCommand("migrate:status", "Show per-plugin migration history in the app binary"), delegateCommand("serve", "Run the app HTTP server"), delegateQueueWorkCommand(), + delegateScheduleRunCommand(), delegateCommand("queue:clear", "Clear pending queued jobs in the app binary"), } } diff --git a/cmd/summer/main_test.go b/cmd/summer/main_test.go index 34fb858..6a3b17b 100644 --- a/cmd/summer/main_test.go +++ b/cmd/summer/main_test.go @@ -19,7 +19,7 @@ func TestToolCommandNames(t *testing.T) { for _, c := range toolCommands() { names = append(names, c.Name) } - for _, want := range []string{"build", "make:plugin", "make:model", "make:migration", "make:command", "make:job", "make:admin-controller", "plugin:add", "dev", "migrate", "migrate:rollback", "migrate:status", "serve", "queue:work", "queue:clear"} { + for _, want := range []string{"build", "make:plugin", "make:model", "make:migration", "make:command", "make:job", "make:admin-controller", "plugin:add", "dev", "migrate", "migrate:rollback", "migrate:status", "serve", "queue:work", "queue:clear", "schedule:run"} { if !slices.Contains(names, want) { t.Fatalf("missing %s in %v", want, names) } @@ -32,6 +32,7 @@ func TestToolCommandNames(t *testing.T) { "make:command": {"[plugin] [name]"}, "make:job": {"[plugin] [name]"}, "make:admin-controller": {"[plugin] [name]"}, + "schedule:run": {"--once"}, } for cmd, wants := range helpWants { var buf bytes.Buffer diff --git a/cmd/summer/runtime.go b/cmd/summer/runtime.go index a916dd6..b3ab945 100644 --- a/cmd/summer/runtime.go +++ b/cmd/summer/runtime.go @@ -6,10 +6,11 @@ import ( "os" "os/exec" "path/filepath" + "strconv" "strings" - "git.golem15.com/golem15/summercms/modules/bonfire" "git.golem15.com/golem15/summercms/internal/build" + "git.golem15.com/golem15/summercms/modules/bonfire" ) func delegateCommand(name, description string) bonfire.Command { @@ -63,6 +64,28 @@ func delegateQueueWorkCommand() bonfire.Command { } } +func delegateScheduleRunCommand() bonfire.Command { + return bonfire.Command{ + Name: "schedule:run", + Description: "Run the scheduler in the app binary (or once for system cron)", + Flags: []bonfire.Flag{{ + Name: "once", + Description: "Run the entries due this minute and exit (for system cron)", + Bare: true, + }}, + Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + argv := []string{"schedule:run"} + if v, ok := in.Flag("once"); ok { + if once, _ := strconv.ParseBool(v); once { + argv = append(argv, "--once") + } + } + argv = append(argv, in.Args()...) + return runAppBinary(ctx, argv, out) + }, + } +} + func runAppBinary(ctx context.Context, argv []string, out bonfire.Output) error { if len(argv) == 0 { return fmt.Errorf("summer: missing app command") diff --git a/modules/conga/README.md b/modules/conga/README.md index fbb81a4..8091a75 100644 --- a/modules/conga/README.md +++ b/modules/conga/README.md @@ -24,7 +24,7 @@ The scheduler is the Go form of WinterCMS `registerSchedule`. Plugins declare re - Workers: `conga.StartWorker` registers every plugin job and starts one River client; `conga.WorkerOptions.Queues` limits it to some queues, and an unknown queue is `conga.ErrUnknownQueue` listing the known ones. `conga.Worker.Stop` stops it gracefully and cancels running jobs when its context ends. `conga.StartServeWorker` is the variant the `serve` command uses: it starts nothing when `queue.work_in_serve` is false. An app without jobs still gets a worker that starts and idles. - Scheduled commands: every worker carries one River periodic job per `pact.HasSchedule` entry, in plugin activation order then declaration order, with the id `[]:`. Only the elected leader enqueues. Each run is a `conga.ScheduledCommandArgs` job on `conga.QueueScheduled` with one attempt (an interrupted run is not retried; the next period runs normally) and unique by args within its cadence period, so a leader failover cannot double-enqueue a period. The worker runs a job only when its entry exists in the compiled schedule and its command and arguments match that entry exactly, so a forged `river_job` row cannot run an arbitrary command. An unregistered command, or an app with no published `bonfire.Catalog`, is logged at Warn (`schedule: command not registered; skipping`) and skipped without failing the worker or other entries. Command output is logged line by line at Info with a `command` attribute; failures are logged with the duration. - Wall-clock schedules: `conga.Daily` fires at the next `hh:mm` in its location and `conga.Every` at the next multiple of its interval since local midnight, so a restart never delays a daily run by up to a day the way `river.PeriodicInterval(24h)` would. On a DST day `conga.Daily` keeps the wall-clock time. A schedule entry with an empty command, a zero cadence, a daily time out of range, or an interval under one second or not dividing 24h fails the worker start with an error naming the plugin id and entry index. -- Commands: `conga.RuntimeCommands` adds `queue:work` and `queue:clear` to the application binary. +- Commands: `conga.RuntimeCommands` adds `queue:work`, `schedule:run` and `queue:clear` to the application binary. `schedule:run` runs a scheduler-only worker in its own process, and `schedule:run --once` runs the entries due in the current minute without River, for system cron. ## Usage @@ -150,7 +150,7 @@ defer w.Stop(context.Background()) | `conga.StartWorker` | Registers plugin jobs and starts a River worker client carrying the plugins' periodic schedule jobs. | | `conga.StartServeWorker` | The worker of the `serve` command; nil when `queue.work_in_serve` is false. | | `conga.WorkerOptions` | Selects the queues a worker runs. | -| `conga.RuntimeCommands` | Returns the `queue:work` and `queue:clear` commands. | +| `conga.RuntimeCommands` | Returns the `queue:work`, `schedule:run` and `queue:clear` commands. | | `conga.Worker` | A running worker; `conga.Worker.Stop` stops it and `conga.Worker.Queues` lists its queues. | | `conga.Daily` | `river.PeriodicSchedule` firing at `Hour:Minute` every day in `Loc` (UTC when nil). | | `conga.Every` | `river.PeriodicSchedule` firing at every multiple of `Interval` since local midnight in `Loc` (UTC when nil). | @@ -185,11 +185,12 @@ queues: ## CLI commands -`conga.RuntimeCommands` adds these commands to the application binary. Both open the database through `lagoon.OpenFromApp`, so they need `database.dsn` and `app.key`. +`conga.RuntimeCommands` adds these commands to the application binary. They open the database through `lagoon.OpenFromApp`, so they need `database.dsn` and `app.key`; `schedule:run --once` opens it only when an entry is due. | Command | Arguments and flags | Description | |---------|---------------------|-------------| | `queue:work` | `--queue `, repeatable | Runs a job worker in the foreground on the named queues (default: every known queue) until SIGINT or SIGTERM, then stops it within 10 seconds. An unknown queue is an error that lists the known ones. | +| `schedule:run` | `--once` (bare) | Without `--once`: runs a worker on the `scheduled` queue that carries every plugin's periodic jobs, prints `scheduler started` and runs until SIGINT or SIGTERM, then stops within 10 seconds. With `--once`: runs, in-process and without River, every entry due in the current minute of `app.timezone` (a daily entry at its hour and minute; `Every(d)` of a minute or less on every run, a longer one when the minutes since midnight are a multiple of `d`), printing `Running scheduled command: ` per entry or `No scheduled commands are ready to run.`. An unregistered command prints a warning and is skipped. The exit status is the first command error, after every due entry ran. There is no overlap lock: two runs in one minute run the due entries twice, as Laravel does. System cron: `* * * * * cd /app && ./bin/app schedule:run --once`. | | `queue:clear` | `[queue]` (default `default`) | Deletes the available, scheduled and retryable jobs of one queue in batches until none are left and prints `Cleared N jobs`. Running jobs are never touched. | ## Dependencies diff --git a/modules/conga/commands.go b/modules/conga/commands.go index b352be1..a0e9ce5 100644 --- a/modules/conga/commands.go +++ b/modules/conga/commands.go @@ -5,6 +5,7 @@ import ( "fmt" "os" "os/signal" + "strconv" "strings" "syscall" "time" @@ -12,16 +13,18 @@ import ( "git.golem15.com/golem15/summercms/modules/backpack" "git.golem15.com/golem15/summercms/modules/bonfire" "git.golem15.com/golem15/summercms/modules/lagoon" + "git.golem15.com/golem15/summercms/modules/pact" "git.golem15.com/golem15/summercms/modules/party" "github.com/riverqueue/river" "github.com/riverqueue/river/rivertype" + "gorm.io/gorm" ) // clearBatch is the number of jobs one JobDeleteMany pass removes. const clearBatch = 10000 -// RuntimeCommands returns the queue:work and queue:clear commands for the -// application binary. +// RuntimeCommands returns the queue:work, schedule:run and queue:clear +// commands for the application binary. func RuntimeCommands(app *backpack.App, plugins []party.Plugin) []bonfire.Command { return []bonfire.Command{ { @@ -38,6 +41,25 @@ func RuntimeCommands(app *backpack.App, plugins []party.Plugin) []bonfire.Comman }) }, }, + { + Name: "schedule:run", + Description: "Run the scheduler in the foreground (or once for system cron)", + Flags: []bonfire.Flag{{ + Name: "once", + Description: "Run the entries due this minute and exit (for system cron)", + Bare: true, + }}, + Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + if v, ok := in.Flag("once"); ok { + if once, _ := strconv.ParseBool(v); once { + return scheduleRunOnce(ctx, app, plugins, scheduleNow(), out) + } + } + return withDB(ctx, app, func() error { + return runScheduler(ctx, app, plugins, out) + }) + }, + }, { Name: "queue:clear", Description: "Clear all queued jobs, by deleting all pending jobs.", @@ -73,6 +95,91 @@ func work(ctx context.Context, app *backpack.App, plugins []party.Plugin, queues return w.Stop(stopCtx) } +// runScheduler runs a scheduler-only worker (the scheduled queue plus every +// plugin's periodic jobs) until SIGINT or SIGTERM. +func runScheduler(ctx context.Context, app *backpack.App, plugins []party.Plugin, out bonfire.Output) error { + ctx, stop := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM) + defer stop() + w, err := StartWorker(ctx, app, plugins, WorkerOptions{Queues: []string{QueueScheduled}}) + if err != nil { + return err + } + out.Info("scheduler started") + <-ctx.Done() + stopCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + return w.Stop(stopCtx) +} + +// scheduleRunOnce runs, without River, every schedule entry due in the +// minute of now in app.timezone (Laravel schedule:run for system cron). An +// unregistered command is warned about and skipped; the returned error is the +// first command error, after every due entry ran. +func scheduleRunOnce(ctx context.Context, app *backpack.App, plugins []party.Plugin, now time.Time, out bonfire.Output) error { + entries, err := scheduleEntries(app, plugins) + if err != nil { + return err + } + loc, err := appLocation(app) + if err != nil { + return err + } + local := now.In(loc) + var due []scheduleEntry + for _, e := range entries { + if dueAt(e.cmd.Cadence, local) { + due = append(due, e) + } + } + if len(due) == 0 { + out.Println("No scheduled commands are ready to run.") + return nil + } + return withAppDB(ctx, app, func() error { + log := loggerFromApp(app) + var first error + for _, e := range due { + if _, ok := scheduledCatalog(app, log, e.cmd.Command); !ok { + out.Warning(fmt.Sprintf("Scheduled command %s is not registered; skipping", e.cmd.Command)) + continue + } + out.Printf("Running scheduled command: %s\n", strings.Join(append([]string{e.cmd.Command}, e.cmd.Args...), " ")) + if _, err := callScheduled(ctx, app, log, e.cmd.Command, e.cmd.Args, out); err != nil && first == nil { + first = err + } + } + return first + }) +} + +// dueAt reports whether a cadence is due in the wall-clock minute of t: a +// daily cadence at its hour and minute, Every(d) of a minute or less on every +// run, and a longer Every(d) when the minutes since midnight are a multiple +// of d. +func dueAt(c pact.Cadence, t time.Time) bool { + if h, m, ok := c.At(); ok { + return t.Hour() == h && t.Minute() == m + } + d := c.Interval() + if d <= 0 { + return false + } + if d <= time.Minute { + return true + } + sinceMidnight := time.Duration(t.Hour()*60+t.Minute()) * time.Minute + return sinceMidnight%d == 0 +} + +// withAppDB runs fn with the app's database: the published one when there is +// one, otherwise one opened for the command the way withDB does. +func withAppDB(ctx context.Context, app *backpack.App, fn func() error) error { + if gdb, ok := app.Lookup[*gorm.DB](); ok && gdb != nil { + return fn() + } + return withDB(ctx, app, fn) +} + // clearQueue deletes the available, scheduled and retryable jobs of queue in // batches until a pass deletes none. River never deletes running jobs. func clearQueue(ctx context.Context, app *backpack.App, queue string, out bonfire.Output) error { diff --git a/modules/conga/schedule_test.go b/modules/conga/schedule_test.go index c3b8303..a12db1a 100644 --- a/modules/conga/schedule_test.go +++ b/modules/conga/schedule_test.go @@ -1,7 +1,9 @@ package conga import ( + "bytes" "context" + "errors" "fmt" "log/slog" "strings" @@ -11,8 +13,10 @@ import ( "git.golem15.com/golem15/summercms/modules/backpack" "git.golem15.com/golem15/summercms/modules/bonfire" + "git.golem15.com/golem15/summercms/modules/compass" "git.golem15.com/golem15/summercms/modules/pact" "git.golem15.com/golem15/summercms/modules/party" + "gorm.io/gorm" ) type schedulePlugin struct { @@ -248,3 +252,356 @@ func TestScheduleRunsCommand(t *testing.T) { logs.waitFor(t, slog.LevelInfo, "tick a", "command", "acme:tick") logs.waitFor(t, slog.LevelWarn, "schedule: command not registered; skipping", "command", "acme:missing") } + +// onceRun runs `schedule:run --once` for plugins at the given app time and +// returns its output. +func onceRun(t *testing.T, app *backpack.App, plugins []party.Plugin, now time.Time) (string, error) { + t.Helper() + prev := scheduleNow + scheduleNow = func() time.Time { return now } + t.Cleanup(func() { scheduleNow = prev }) + var out bytes.Buffer + root, err := bonfire.NewRoot("acme", RuntimeCommands(app, plugins), &out) + if err != nil { + t.Fatal(err) + } + root.SetArgs([]string{"schedule:run", "--once"}) + err = root.ExecuteContext(t.Context()) + return out.String(), err +} + +// onceApp is an app with a published catalog of acme:tick (records its +// args) and acme:boom (fails), and a published *gorm.DB so --once does not +// open a connection. +func onceApp(t *testing.T, extra map[string]any) (*backpack.App, *[][]string, *captureHandler) { + t.Helper() + cfg, err := compass.Open(compass.Options{Dir: t.TempDir(), Env: "testing", Environ: []string{}}) + if err != nil { + t.Fatal(err) + } + for k, v := range extra { + if err := cfg.Set(k, v); err != nil { + t.Fatal(err) + } + } + app := backpack.New(cfg) + logs := &captureHandler{} + if err := app.Publish(slog.New(logs)); err != nil { + t.Fatal(err) + } + if err := app.Publish(&gorm.DB{}); err != nil { + t.Fatal(err) + } + var ran [][]string + catalog := bonfire.NewCatalog([]bonfire.Command{ + {Name: "acme:tick", Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + ran = append(ran, append([]string{"acme:tick"}, in.Args()...)) + return nil + }}, + {Name: "acme:boom", Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + ran = append(ran, []string{"acme:boom"}) + return errors.New("boom") + }}, + }) + if err := app.Publish(catalog); err != nil { + t.Fatal(err) + } + return app, &ran, logs +} + +// TestScheduleRunOnce covers the Laravel schedule:run minute-match semantics +// of `schedule:run --once` (D-18, CLI-04). +func TestScheduleRunOnce(t *testing.T) { + utc := time.UTC + + t.Run("daily_due_at_midnight_only", func(t *testing.T) { + app, ran, _ := onceApp(t, nil) + plugins := schedulePlugins(pact.ScheduledCommand{Command: "acme:tick", Cadence: pact.Daily()}) + out, err := onceRun(t, app, plugins, mustTime(t, utc, "2026-03-10 00:00:42")) + if err != nil || !strings.Contains(out, "Running scheduled command: acme:tick\n") || len(*ran) != 1 { + t.Fatalf("00:00: out %q, ran %v, err %v", out, *ran, err) + } + out, err = onceRun(t, app, plugins, mustTime(t, utc, "2026-03-10 00:01:00")) + if err != nil || !strings.Contains(out, "No scheduled commands are ready to run.") || len(*ran) != 1 { + t.Fatalf("00:01: out %q, ran %v, err %v", out, *ran, err) + } + }) + + t.Run("every_five_minutes", func(t *testing.T) { + app, ran, _ := onceApp(t, nil) + plugins := schedulePlugins(pact.ScheduledCommand{Command: "acme:tick", Args: []string{"x", "y"}, Cadence: pact.Every(5 * time.Minute)}) + out, err := onceRun(t, app, plugins, mustTime(t, utc, "2026-03-10 10:05:00")) + if err != nil || !strings.Contains(out, "Running scheduled command: acme:tick x y\n") { + t.Fatalf("10:05: out %q, err %v", out, err) + } + if fmt.Sprint(*ran) != "[[acme:tick x y]]" { + t.Fatalf("ran = %v", *ran) + } + out, err = onceRun(t, app, plugins, mustTime(t, utc, "2026-03-10 10:07:00")) + if err != nil || !strings.Contains(out, "No scheduled commands are ready to run.") || len(*ran) != 1 { + t.Fatalf("10:07: out %q, ran %v, err %v", out, *ran, err) + } + }) + + t.Run("app_timezone", func(t *testing.T) { + app, ran, _ := onceApp(t, map[string]any{"app.timezone": "Europe/Warsaw"}) + plugins := schedulePlugins(pact.ScheduledCommand{Command: "acme:tick", Cadence: pact.DailyAt(1, 0)}) + // 00:00 UTC in January is 01:00 in Warsaw. + if _, err := onceRun(t, app, plugins, mustTime(t, utc, "2026-01-10 00:00:00")); err != nil || len(*ran) != 1 { + t.Fatalf("ran %v, err %v", *ran, err) + } + }) + + t.Run("unregistered_command_warns_and_others_run", func(t *testing.T) { + app, ran, logs := onceApp(t, nil) + plugins := schedulePlugins( + pact.ScheduledCommand{Command: "acme:missing", Cadence: pact.Daily()}, + pact.ScheduledCommand{Command: "acme:tick", Cadence: pact.Daily()}, + ) + out, err := onceRun(t, app, plugins, mustTime(t, utc, "2026-03-10 00:00:00")) + if err != nil { + t.Fatalf("err = %v, want exit 0", err) + } + if !strings.Contains(out, "acme:missing") || !strings.Contains(out, "not registered") { + t.Fatalf("no warning naming acme:missing: %q", out) + } + if fmt.Sprint(*ran) != "[[acme:tick]]" { + t.Fatalf("ran = %v, want acme:tick to still run", *ran) + } + if !logs.find(slog.LevelWarn, "schedule: command not registered; skipping", "command", "acme:missing") { + t.Fatal("no Warn log for acme:missing") + } + }) + + t.Run("first_error_after_all_due_ran", func(t *testing.T) { + app, ran, _ := onceApp(t, nil) + plugins := schedulePlugins( + pact.ScheduledCommand{Command: "acme:boom", Cadence: pact.Every(time.Minute)}, + pact.ScheduledCommand{Command: "acme:tick", Cadence: pact.Every(30 * time.Second)}, + ) + _, err := onceRun(t, app, plugins, mustTime(t, utc, "2026-03-10 10:07:00")) + if err == nil || !strings.Contains(err.Error(), "boom") { + t.Fatalf("err = %v, want the acme:boom error", err) + } + if fmt.Sprint(*ran) != "[[acme:boom] [acme:tick]]" { + t.Fatalf("ran = %v, want both entries", *ran) + } + }) + + t.Run("invalid_entry_fails", func(t *testing.T) { + app, _, _ := onceApp(t, nil) + _, err := onceRun(t, app, schedulePlugins(pact.ScheduledCommand{Command: "acme:tick"}), mustTime(t, utc, "2026-03-10 00:00:00")) + if err == nil || !strings.Contains(err.Error(), "plugin acme.test schedule entry 0") { + t.Fatalf("err = %v", err) + } + }) +} + +// TestScheduledEntryMismatchSkipped covers T-11-09: a scheduled command job +// whose command or args differ from its compiled entry, or whose entry does +// not exist, is logged and skipped without running anything, including when +// a forged row reaches a running worker. +func TestScheduledEntryMismatchSkipped(t *testing.T) { + db, dsn := migratedDB(t) + app, _ := testApp(t, db, dsn, nil) + logs := &captureHandler{} + if err := app.Publish(slog.New(logs)); err != nil { + t.Fatal(err) + } + var mu sync.Mutex + var ran [][]string + if err := app.Publish(bonfire.NewCatalog([]bonfire.Command{ + {Name: "acme:tick", Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + mu.Lock() + ran = append(ran, append([]string{"acme:tick"}, in.Args()...)) + mu.Unlock() + return nil + }}, + {Name: "acme:other", Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + mu.Lock() + ran = append(ran, []string{"acme:other"}) + mu.Unlock() + return nil + }}, + })); err != nil { + t.Fatal(err) + } + w, err := StartWorker(t.Context(), app, schedulePlugins( + pact.ScheduledCommand{Command: "acme:tick", Args: []string{"safe"}, Cadence: pact.Daily()}, + ), WorkerOptions{}) + if err != nil { + t.Fatal(err) + } + defer func() { _ = w.Stop(context.Background()) }() + m, err := From(app) + if err != nil { + t.Fatal(err) + } + const entry = "acme.test[0]:acme:tick" + const skip = "schedule: job does not match a compiled schedule entry; skipping" + + for _, forged := range []ScheduledCommandArgs{ + {Entry: entry, Command: "acme:other", Args: []string{"safe"}}, + {Entry: entry, Command: "acme:tick", Args: []string{"--evil"}}, + {Entry: entry, Command: "acme:tick"}, + {Entry: "acme.test[9]:acme:tick", Command: "acme:tick", Args: []string{"safe"}}, + } { + if err := m.runScheduled(t.Context(), forged); err != nil { + t.Fatalf("%+v: err %v", forged, err) + } + if !logs.find(slog.LevelWarn, skip, "command", forged.Command) { + t.Fatalf("%+v: no skip warning", forged) + } + } + + // A forged river_job row on the scheduled queue reaches the worker and + // is skipped the same way. + if err := m.Enqueue(t.Context(), nil, ScheduledCommandArgs{Entry: entry, Command: "acme:other"}, EnqueueOpts{Queue: QueueScheduled}); err != nil { + t.Fatal(err) + } + deadline := time.Now().Add(10 * time.Second) + for { + var n int + if err := db.QueryRowContext(t.Context(), `SELECT count(*) FROM river_job WHERE kind = 'summer.scheduled_command' AND state = 'completed'`).Scan(&n); err != nil { + t.Fatal(err) + } + if n == 1 { + break + } + if time.Now().After(deadline) { + t.Fatal("forged job was not worked within 10s") + } + time.Sleep(25 * time.Millisecond) + } + snapshot := func() string { + mu.Lock() + defer mu.Unlock() + return fmt.Sprint(ran) + } + if got := snapshot(); got != "[]" { + t.Fatalf("forged jobs ran commands: %s", got) + } + + // The matching entry does run. + if err := m.runScheduled(t.Context(), ScheduledCommandArgs{Entry: entry, Command: "acme:tick", Args: []string{"safe"}}); err != nil { + t.Fatal(err) + } + if got := snapshot(); got != "[[acme:tick safe]]" { + t.Fatalf("ran = %s", got) + } +} + +// TestScheduleUniqueByPeriod covers T-11-17: two inserts of one periodic +// entry inside its cadence period leave one river_job row, while another +// entry in the same period gets its own row. +func TestScheduleUniqueByPeriod(t *testing.T) { + db, dsn := migratedDB(t) + app, _ := testApp(t, db, dsn, nil) + entries, err := scheduleEntries(app, schedulePlugins( + pact.ScheduledCommand{Command: "acme:tick", Cadence: pact.Daily()}, + pact.ScheduledCommand{Command: "acme:tock", Cadence: pact.Daily()}, + )) + if err != nil { + t.Fatal(err) + } + m, err := From(app) + if err != nil { + t.Fatal(err) + } + client, err := m.insertClient() + if err != nil { + t.Fatal(err) + } + count := func() int { + var n int + if err := db.QueryRowContext(t.Context(), `SELECT count(*) FROM river_job WHERE kind = 'summer.scheduled_command'`).Scan(&n); err != nil { + t.Fatal(err) + } + return n + } + for i := range 2 { + args, opts := entries[0].construct() + res, err := client.Insert(t.Context(), args, opts) + if err != nil { + t.Fatal(err) + } + if i == 1 && !res.UniqueSkippedAsDuplicate { + t.Fatal("second insert in the period was not skipped as a duplicate") + } + if opts.Queue != QueueScheduled || opts.MaxAttempts != 1 || opts.UniqueOpts.ByPeriod != 24*time.Hour { + t.Fatalf("insert opts = %+v", opts) + } + } + if n := count(); n != 1 { + t.Fatalf("river_job rows = %d, want 1", n) + } + args, opts := entries[1].construct() + if _, err := client.Insert(t.Context(), args, opts); err != nil { + t.Fatal(err) + } + if n := count(); n != 2 { + t.Fatalf("river_job rows = %d, want 2 (one per entry)", n) + } +} + +// TestScheduleRunForeground covers D-18: schedule:run without --once runs a +// scheduler-only worker (scheduled queue plus periodic jobs) until its +// context ends. +func TestScheduleRunForeground(t *testing.T) { + _, dsn := migratedDB(t) + app := commandApp(t, dsn, nil) + ticked := make(chan struct{}, 16) + if err := app.Publish(bonfire.NewCatalog([]bonfire.Command{{ + Name: "acme:tick", + Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + select { + case ticked <- struct{}{}: + default: + } + return nil + }, + }})); err != nil { + t.Fatal(err) + } + out := newSyncBuffer() + root, err := bonfire.NewRoot("acme", RuntimeCommands(app, schedulePlugins( + pact.ScheduledCommand{Command: "acme:tick", Cadence: pact.Every(time.Second)}, + )), out) + if err != nil { + t.Fatal(err) + } + root.SetArgs([]string{"schedule:run"}) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + done := make(chan error, 1) + go func() { done <- root.ExecuteContext(ctx) }() + select { + case <-ticked: + case err := <-done: + t.Fatalf("schedule:run returned early: %v (%q)", err, out.String()) + case <-time.After(10 * time.Second): + t.Fatalf("scheduled command did not run: %q", out.String()) + } + if !strings.Contains(out.String(), "scheduler started") { + t.Fatalf("out = %q", out.String()) + } + m, err := From(app) + if err != nil { + t.Fatal(err) + } + m.mu.Lock() + running := m.worker != nil + m.mu.Unlock() + if !running { + t.Fatal("no worker running") + } + cancel() + select { + case err := <-done: + if err != nil { + t.Fatalf("schedule:run returned %v", err) + } + case <-time.After(15 * time.Second): + t.Fatal("schedule:run did not stop") + } +} diff --git a/modules/conga/scheduler.go b/modules/conga/scheduler.go index a16cb8d..6b7c9ff 100644 --- a/modules/conga/scheduler.go +++ b/modules/conga/scheduler.go @@ -145,16 +145,8 @@ func (m *Manager) runScheduled(ctx context.Context, a ScheduledCommandArgs) erro // at Warn and skipped (skipped true, nil error); a command error is logged // with the duration and returned. func callScheduled(ctx context.Context, app *backpack.App, log *slog.Logger, command string, args []string, out io.Writer) (bool, error) { - var cat *bonfire.Catalog - if app != nil { - cat, _ = app.Lookup[*bonfire.Catalog]() - } - if cat == nil { - log.Warn("schedule: no command catalog published; skipping", "command", command) - return true, nil - } - if !cat.Has(command) { - log.Warn("schedule: command not registered; skipping", "command", command) + cat, ok := scheduledCatalog(app, log, command) + if !ok { return true, nil } start := time.Now() @@ -168,6 +160,24 @@ func callScheduled(ctx context.Context, app *backpack.App, log *slog.Logger, com return false, nil } +// scheduledCatalog returns the app's catalog when it holds command. A missing +// catalog or command is logged at Warn and reported as false. +func scheduledCatalog(app *backpack.App, log *slog.Logger, command string) (*bonfire.Catalog, bool) { + var cat *bonfire.Catalog + if app != nil { + cat, _ = app.Lookup[*bonfire.Catalog]() + } + if cat == nil { + log.Warn("schedule: no command catalog published; skipping", "command", command) + return nil, false + } + if !cat.Has(command) { + log.Warn("schedule: command not registered; skipping", "command", command) + return nil, false + } + return cat, true +} + // logWriter logs each output line of a scheduled command at Info. type logWriter struct { log *slog.Logger @@ -214,3 +224,6 @@ func (w *logWriter) emit(line string) { } w.log.Info(line, "command", w.command) } + +// scheduleNow is the clock of schedule:run --once (tests replace it). +var scheduleNow = time.Now