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
This commit is contained in:
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user