- deferred:purge [--days] in lagoon.RuntimeCommands (purge_days, default 5)
- lagoon.FrameworkSchedule entry at purge_at (default 03:00, empty disables)
- conga prepends framework entries as summercms.lagoon[i]:<command>
- pact.Relation{Before,After}{Create,Update,Delete} optional hooks
- lagoon, conga and pact READMEs, scheduling and setup docs
846 lines
31 KiB
Go
846 lines
31 KiB
Go
package conga
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"strings"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"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 {
|
|
id string
|
|
schedule []pact.ScheduledCommand
|
|
}
|
|
|
|
func (p *schedulePlugin) ID() string { return p.id }
|
|
func (p *schedulePlugin) Requires() []string { return nil }
|
|
func (p *schedulePlugin) Register(*backpack.App) error { return nil }
|
|
func (p *schedulePlugin) Boot(*backpack.App) error { return nil }
|
|
func (p *schedulePlugin) Schedule() []pact.ScheduledCommand { return p.schedule }
|
|
func schedulePlugins(entries ...pact.ScheduledCommand) []party.Plugin {
|
|
return []party.Plugin{&schedulePlugin{id: "acme.test", schedule: entries}}
|
|
}
|
|
|
|
// captureHandler records log messages with their string attributes.
|
|
type captureHandler struct {
|
|
mu sync.Mutex
|
|
records []capturedRecord
|
|
}
|
|
|
|
type capturedRecord struct {
|
|
level slog.Level
|
|
msg string
|
|
attrs map[string]string
|
|
}
|
|
|
|
func (h *captureHandler) Enabled(context.Context, slog.Level) bool { return true }
|
|
func (h *captureHandler) WithAttrs([]slog.Attr) slog.Handler { return h }
|
|
func (h *captureHandler) WithGroup(string) slog.Handler { return h }
|
|
func (h *captureHandler) Handle(_ context.Context, r slog.Record) error {
|
|
rec := capturedRecord{level: r.Level, msg: r.Message, attrs: map[string]string{}}
|
|
r.Attrs(func(a slog.Attr) bool {
|
|
rec.attrs[a.Key] = a.Value.String()
|
|
return true
|
|
})
|
|
h.mu.Lock()
|
|
h.records = append(h.records, rec)
|
|
h.mu.Unlock()
|
|
return nil
|
|
}
|
|
|
|
// find returns the first record with msg whose attrs contain key=value.
|
|
func (h *captureHandler) find(level slog.Level, msg, key, value string) bool {
|
|
h.mu.Lock()
|
|
defer h.mu.Unlock()
|
|
for _, r := range h.records {
|
|
if r.level == level && r.msg == msg && r.attrs[key] == value {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
func (h *captureHandler) waitFor(t *testing.T, level slog.Level, msg, key, value string) {
|
|
t.Helper()
|
|
deadline := time.Now().Add(10 * time.Second)
|
|
for !h.find(level, msg, key, value) {
|
|
if time.Now().After(deadline) {
|
|
t.Fatalf("no %s log %q with %s=%s", level, msg, key, value)
|
|
}
|
|
time.Sleep(25 * time.Millisecond)
|
|
}
|
|
}
|
|
|
|
func mustTime(t *testing.T, loc *time.Location, s string) time.Time {
|
|
t.Helper()
|
|
tm, err := time.ParseInLocation("2006-01-02 15:04:05", s, loc)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return tm
|
|
}
|
|
|
|
// TestScheduleNext covers the CLI-04 adjacency edges and wall-clock DST
|
|
// behaviour of the Daily and Every schedules.
|
|
func TestScheduleNext(t *testing.T) {
|
|
utc := time.UTC
|
|
warsaw, err := time.LoadLocation("Europe/Warsaw")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
cases := []struct {
|
|
name string
|
|
s interface{ Next(time.Time) time.Time }
|
|
now time.Time
|
|
want time.Time
|
|
}{
|
|
{"daily_exactly_on_boundary_is_next_day", Daily{Loc: utc}, mustTime(t, utc, "2026-03-10 00:00:00"), mustTime(t, utc, "2026-03-11 00:00:00")},
|
|
{"daily_just_before", Daily{Loc: utc}, mustTime(t, utc, "2026-03-10 23:59:59"), mustTime(t, utc, "2026-03-11 00:00:00")},
|
|
{"daily_at_later_today", Daily{Hour: 3, Minute: 30, Loc: utc}, mustTime(t, utc, "2026-03-10 03:29:59"), mustTime(t, utc, "2026-03-10 03:30:00")},
|
|
{"daily_nil_loc_is_utc", Daily{Hour: 1}, mustTime(t, utc, "2026-03-10 01:00:00"), mustTime(t, utc, "2026-03-11 01:00:00")},
|
|
{"daily_converts_input_location", Daily{Loc: warsaw}, mustTime(t, utc, "2026-01-10 22:30:00"), mustTime(t, warsaw, "2026-01-11 00:00:00")},
|
|
{"daily_spring_forward_keeps_wall_clock", Daily{Hour: 12, Loc: warsaw}, mustTime(t, warsaw, "2026-03-28 12:00:00"), mustTime(t, warsaw, "2026-03-29 12:00:00")},
|
|
{"daily_fall_back_keeps_wall_clock", Daily{Hour: 12, Loc: warsaw}, mustTime(t, warsaw, "2026-10-24 12:00:00"), mustTime(t, warsaw, "2026-10-25 12:00:00")},
|
|
{"every_exactly_on_multiple_is_strictly_after", Every{Interval: 5 * time.Minute, Loc: utc}, mustTime(t, utc, "2026-03-10 10:05:00"), mustTime(t, utc, "2026-03-10 10:10:00")},
|
|
{"every_between_multiples", Every{Interval: 5 * time.Minute, Loc: utc}, mustTime(t, utc, "2026-03-10 10:07:13"), mustTime(t, utc, "2026-03-10 10:10:00")},
|
|
{"every_rolls_to_midnight", Every{Interval: 5 * time.Minute, Loc: utc}, mustTime(t, utc, "2026-03-10 23:57:00"), mustTime(t, utc, "2026-03-11 00:00:00")},
|
|
{"every_second_on_boundary", Every{Interval: time.Second, Loc: utc}, mustTime(t, utc, "2026-03-10 10:00:00"), mustTime(t, utc, "2026-03-10 10:00:01")},
|
|
{"every_hour_in_location", Every{Interval: time.Hour, Loc: warsaw}, mustTime(t, warsaw, "2026-01-10 09:15:00"), mustTime(t, warsaw, "2026-01-10 10:00:00")},
|
|
}
|
|
for _, c := range cases {
|
|
t.Run(c.name, func(t *testing.T) {
|
|
got := c.s.Next(c.now)
|
|
if !got.Equal(c.want) {
|
|
t.Fatalf("Next(%s) = %s, want %s", c.now, got, c.want)
|
|
}
|
|
})
|
|
}
|
|
|
|
// The DST-day runs are 23h and 25h apart but stay at 12:00 local.
|
|
spring := Daily{Hour: 12, Loc: warsaw}.Next(mustTime(t, warsaw, "2026-03-28 12:00:00"))
|
|
if d := spring.Sub(mustTime(t, warsaw, "2026-03-28 12:00:00")); d != 23*time.Hour {
|
|
t.Fatalf("spring-forward gap = %s, want 23h", d)
|
|
}
|
|
if h, m, _ := spring.Clock(); h != 12 || m != 0 {
|
|
t.Fatalf("spring-forward wall clock = %02d:%02d", h, m)
|
|
}
|
|
|
|
for _, bad := range []pact.Cadence{{}, pact.Every(0), pact.Every(500 * time.Millisecond), pact.Every(7 * time.Minute), pact.DailyAt(24, 0), pact.DailyAt(0, 60)} {
|
|
if _, _, err := scheduleFor(bad, utc); err == nil {
|
|
t.Fatalf("scheduleFor(%+v) accepted an invalid cadence", bad)
|
|
}
|
|
}
|
|
if _, period, err := scheduleFor(pact.Daily(), utc); err != nil || period != 24*time.Hour {
|
|
t.Fatalf("daily period = %s, err %v", period, err)
|
|
}
|
|
if _, period, err := scheduleFor(pact.Every(15*time.Minute), utc); err != nil || period != 15*time.Minute {
|
|
t.Fatalf("every period = %s, err %v", period, err)
|
|
}
|
|
}
|
|
|
|
// TestScheduleEntries covers the CLI-04 empty, invalid and ordering edges.
|
|
func TestScheduleEntries(t *testing.T) {
|
|
app := backpack.New(nil)
|
|
jobs, table, err := periodicJobs(app, []party.Plugin{
|
|
&schedulePlugin{id: "acme.none"},
|
|
&schedulePlugin{id: "acme.empty", schedule: []pact.ScheduledCommand{}},
|
|
})
|
|
if err != nil || len(jobs) != 0 || len(table) != 0 {
|
|
t.Fatalf("empty schedules: jobs %d, table %d, err %v", len(jobs), len(table), err)
|
|
}
|
|
|
|
entries, err := scheduleEntries(app, []party.Plugin{
|
|
&schedulePlugin{id: "acme.b", schedule: []pact.ScheduledCommand{
|
|
{Command: "acme:one", Cadence: pact.Daily()},
|
|
{Command: "acme:two", Cadence: pact.Daily()},
|
|
}},
|
|
&schedulePlugin{id: "acme.a", schedule: []pact.ScheduledCommand{{Command: "acme:three", Cadence: pact.Every(time.Minute)}}},
|
|
})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var ids []string
|
|
for _, e := range entries {
|
|
ids = append(ids, e.id)
|
|
}
|
|
if got := strings.Join(ids, ","); got != "acme.b[0]:acme:one,acme.b[1]:acme:two,acme.a[0]:acme:three" {
|
|
t.Fatalf("entry ids = %s", got)
|
|
}
|
|
|
|
for _, c := range []struct {
|
|
entry pact.ScheduledCommand
|
|
want string
|
|
}{
|
|
{pact.ScheduledCommand{Command: " ", Cadence: pact.Daily()}, "plugin acme.bad schedule entry 1"},
|
|
{pact.ScheduledCommand{Command: "acme:x"}, "plugin acme.bad schedule entry 1"},
|
|
{pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Every(7 * time.Minute)}, "plugin acme.bad schedule entry 1"},
|
|
} {
|
|
_, _, err := periodicJobs(app, []party.Plugin{&schedulePlugin{id: "acme.bad", schedule: []pact.ScheduledCommand{
|
|
{Command: "acme:ok", Cadence: pact.Daily()}, c.entry,
|
|
}}})
|
|
if err == nil || !strings.Contains(err.Error(), c.want) {
|
|
t.Fatalf("entry %+v: err = %v, want %q", c.entry, err, c.want)
|
|
}
|
|
}
|
|
|
|
if _, err := appLocation(backpack.New(nil)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
// TestScheduleRunsCommand covers CLI-04 end to end: an Every(1s) entry runs
|
|
// through a River periodic job and bonfire.Call inside a worker, and an
|
|
// unregistered command is skipped with a Warn log (user decision 5).
|
|
func TestScheduleRunsCommand(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)
|
|
}
|
|
ticked := make(chan []string, 16)
|
|
catalog := bonfire.NewCatalog([]bonfire.Command{{
|
|
Name: "acme:tick",
|
|
Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error {
|
|
out.Printf("tick %s\n", strings.Join(in.Args(), " "))
|
|
select {
|
|
case ticked <- in.Args():
|
|
default:
|
|
}
|
|
return nil
|
|
},
|
|
}})
|
|
if err := app.Publish(catalog); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
w, err := StartWorker(t.Context(), app, schedulePlugins(
|
|
pact.ScheduledCommand{Command: "acme:tick", Args: []string{"a"}, Cadence: pact.Every(time.Second)},
|
|
pact.ScheduledCommand{Command: "acme:missing", Cadence: pact.Every(time.Second)},
|
|
), WorkerOptions{})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() {
|
|
if err := w.Stop(context.Background()); err != nil {
|
|
t.Error(err)
|
|
}
|
|
}()
|
|
if got := strings.Join(w.Queues(), ","); got != "default,scheduled" {
|
|
t.Fatalf("queues = %s", got)
|
|
}
|
|
|
|
select {
|
|
case args := <-ticked:
|
|
if fmt.Sprint(args) != "[a]" {
|
|
t.Fatalf("acme:tick args = %v", args)
|
|
}
|
|
case <-time.After(10 * time.Second):
|
|
t.Fatal("acme:tick did not run within 10s")
|
|
}
|
|
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")
|
|
}
|
|
}
|
|
|
|
// TestScheduleValidation covers the CLI-04 refusals: an empty command, a
|
|
// zero cadence, an interval that does not divide 24h and an invalid
|
|
// app.timezone each fail the compile, naming the plugin and entry index,
|
|
// and a worker does not start with them.
|
|
func TestScheduleValidation(t *testing.T) {
|
|
cases := []struct {
|
|
name string
|
|
app *backpack.App
|
|
entry pact.ScheduledCommand
|
|
want string
|
|
}{
|
|
{"empty_command", backpack.New(nil), pact.ScheduledCommand{Command: "", Cadence: pact.Daily()}, "plugin acme.v schedule entry 1: command is empty"},
|
|
{"zero_cadence", backpack.New(nil), pact.ScheduledCommand{Command: "acme:x"}, "plugin acme.v schedule entry 1 (acme:x): cadence is zero"},
|
|
{"non_dividing_interval", backpack.New(nil), pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Every(7 * time.Minute)}, "plugin acme.v schedule entry 1 (acme:x): interval 7m0s does not divide 24h evenly"},
|
|
{"sub_second_interval", backpack.New(nil), pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Every(time.Millisecond)}, "shorter than one second"},
|
|
{"negative_interval", backpack.New(nil), pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Every(-time.Minute)}, "not positive"},
|
|
{"daily_out_of_range", backpack.New(nil), pact.ScheduledCommand{Command: "acme:x", Cadence: pact.DailyAt(12, 75)}, "daily time 12:75 is out of range"},
|
|
}
|
|
for _, c := range cases {
|
|
t.Run(c.name, func(t *testing.T) {
|
|
plugins := []party.Plugin{&schedulePlugin{id: "acme.v", schedule: []pact.ScheduledCommand{{Command: "acme:ok", Cadence: pact.Daily()}, c.entry}}}
|
|
_, err := scheduleEntries(c.app, plugins)
|
|
if err == nil || !strings.Contains(err.Error(), c.want) {
|
|
t.Fatalf("err = %v, want %q", err, c.want)
|
|
}
|
|
if _, err := StartWorker(t.Context(), c.app, plugins, WorkerOptions{}); err == nil {
|
|
t.Fatal("worker started with an invalid schedule")
|
|
}
|
|
})
|
|
}
|
|
t.Run("invalid_timezone", func(t *testing.T) {
|
|
cfg, err := compass.Open(compass.Options{Dir: t.TempDir(), Env: "testing", Environ: []string{}})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := cfg.Set("app.timezone", "Mars/Olympus"); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
app := backpack.New(cfg)
|
|
_, err = scheduleEntries(app, schedulePlugins(pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Daily()}))
|
|
if err == nil || !strings.Contains(err.Error(), `app.timezone "Mars/Olympus"`) {
|
|
t.Fatalf("err = %v", err)
|
|
}
|
|
if _, err := onceRun(t, app, schedulePlugins(pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Daily()}), time.Now()); err == nil {
|
|
t.Fatal("--once ran with an invalid app.timezone")
|
|
}
|
|
})
|
|
}
|
|
|
|
// TestScheduleOrdering covers entry order and ids: plugins in activation
|
|
// order, entries in declaration order, plugins without a schedule skipped,
|
|
// and a compiled table keyed by id with cloned args.
|
|
func TestScheduleOrdering(t *testing.T) {
|
|
args := []string{"--force"}
|
|
plugins := []party.Plugin{
|
|
&schedulePlugin{id: "acme.z", schedule: []pact.ScheduledCommand{
|
|
{Command: "acme:first", Args: args, Cadence: pact.Every(time.Hour)},
|
|
{Command: "acme:second", Cadence: pact.DailyAt(3, 15)},
|
|
}},
|
|
&jobsPlugin{id: "acme.jobs-only"},
|
|
&schedulePlugin{id: "acme.a", schedule: []pact.ScheduledCommand{{Command: "acme:first", Cadence: pact.Daily()}}},
|
|
}
|
|
jobs, table, err := periodicJobs(backpack.New(nil), plugins)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(jobs) != 3 || len(table) != 3 {
|
|
t.Fatalf("jobs %d, table %d, want 3", len(jobs), len(table))
|
|
}
|
|
entries, err := scheduleEntries(backpack.New(nil), plugins)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var ids []string
|
|
for _, e := range entries {
|
|
ids = append(ids, e.id)
|
|
}
|
|
want := []string{"acme.z[0]:acme:first", "acme.z[1]:acme:second", "acme.a[0]:acme:first"}
|
|
if strings.Join(ids, " ") != strings.Join(want, " ") {
|
|
t.Fatalf("ids = %v, want %v", ids, want)
|
|
}
|
|
args[0] = "--mutated"
|
|
if got := table["acme.z[0]:acme:first"].Args; len(got) != 1 || got[0] != "--force" {
|
|
t.Fatalf("compiled args = %v, want a copy of the declared args", got)
|
|
}
|
|
a, opts := entries[1].construct()
|
|
sa, ok := a.(ScheduledCommandArgs)
|
|
if !ok || sa.Entry != "acme.z[1]:acme:second" || sa.Kind() != "summer.scheduled_command" {
|
|
t.Fatalf("constructed args = %#v", a)
|
|
}
|
|
if opts.Queue != QueueScheduled || opts.MaxAttempts != 1 || !opts.UniqueOpts.ByArgs || opts.UniqueOpts.ByPeriod != 24*time.Hour {
|
|
t.Fatalf("insert opts = %+v", opts)
|
|
}
|
|
}
|
|
|
|
// TestScheduleMissingCatalog covers user decision 5 without a published
|
|
// command catalog: a scheduled run is logged at Warn and skipped, and a
|
|
// matching entry whose command fails returns the error.
|
|
func TestScheduleMissingCatalog(t *testing.T) {
|
|
app := backpack.New(nil)
|
|
logs := &captureHandler{}
|
|
if err := app.Publish(slog.New(logs)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
m, err := From(app)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
m.schedule = map[string]pact.ScheduledCommand{"acme.test[0]:acme:tick": {Command: "acme:tick", Cadence: pact.Daily()}}
|
|
if err := m.runScheduled(t.Context(), ScheduledCommandArgs{Entry: "acme.test[0]:acme:tick", Command: "acme:tick"}); err != nil {
|
|
t.Fatalf("run without a catalog = %v, want nil", err)
|
|
}
|
|
if !logs.find(slog.LevelWarn, "schedule: no command catalog published; skipping", "command", "acme:tick") {
|
|
t.Fatal("no Warn log for the missing catalog")
|
|
}
|
|
if err := app.Publish(bonfire.NewCatalog([]bonfire.Command{{Name: "acme:tick", Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error {
|
|
out.Println("partial output")
|
|
return errors.New("tick failed")
|
|
}}})); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
err = m.runScheduled(t.Context(), ScheduledCommandArgs{Entry: "acme.test[0]:acme:tick", Command: "acme:tick"})
|
|
if err == nil || !strings.Contains(err.Error(), "tick failed") {
|
|
t.Fatalf("failing command = %v", err)
|
|
}
|
|
if !logs.find(slog.LevelInfo, "partial output", "command", "acme:tick") {
|
|
t.Fatal("command output was not logged")
|
|
}
|
|
if !logs.find(slog.LevelError, "schedule: command failed", "command", "acme:tick") {
|
|
t.Fatal("no Error log for the failed command")
|
|
}
|
|
}
|
|
|
|
// TestScheduleLogWriter covers the scheduled command output logger: lines
|
|
// split across writes, CRLF, blank lines and a trailing partial line.
|
|
func TestScheduleLogWriter(t *testing.T) {
|
|
logs := &captureHandler{}
|
|
w := newLogWriter(slog.New(logs), "acme:tick")
|
|
for _, chunk := range []string{"first li", "ne\r\n\nsecond\nthi", "rd"} {
|
|
if n, err := w.Write([]byte(chunk)); err != nil || n != len(chunk) {
|
|
t.Fatalf("Write(%q) = %d, %v", chunk, n, err)
|
|
}
|
|
}
|
|
w.Flush()
|
|
w.Flush()
|
|
var msgs []string
|
|
for _, r := range logs.records {
|
|
if r.attrs["command"] != "acme:tick" {
|
|
t.Fatalf("record without the command attribute: %+v", r)
|
|
}
|
|
msgs = append(msgs, r.msg)
|
|
}
|
|
if strings.Join(msgs, "|") != "first line|second|third" {
|
|
t.Fatalf("logged lines = %q", msgs)
|
|
}
|
|
}
|
|
|
|
// TestScheduleDueAt covers the --once minute match, including an Every(d)
|
|
// longer than a minute that does not divide an hour.
|
|
func TestScheduleDueAt(t *testing.T) {
|
|
at := func(hm string) time.Time { return mustTime(t, time.UTC, "2026-03-10 "+hm+":00") }
|
|
cases := []struct {
|
|
name string
|
|
c pact.Cadence
|
|
t time.Time
|
|
want bool
|
|
}{
|
|
{"daily_at_match", pact.DailyAt(3, 15), at("03:15"), true},
|
|
{"daily_at_other_minute", pact.DailyAt(3, 15), at("03:16"), false},
|
|
{"every_minute_always", pact.Every(time.Minute), at("07:13"), true},
|
|
{"every_30s_always", pact.Every(30 * time.Second), at("07:13"), true},
|
|
{"every_90s_at_3m", pact.Every(90 * time.Second), at("00:03"), true},
|
|
{"every_90s_at_1m", pact.Every(90 * time.Second), at("00:01"), false},
|
|
{"every_90s_at_1h", pact.Every(90 * time.Second), at("01:00"), true},
|
|
{"every_hour_on_the_hour", pact.Every(time.Hour), at("13:00"), true},
|
|
{"every_hour_off_the_hour", pact.Every(time.Hour), at("13:30"), false},
|
|
{"zero_cadence", pact.Cadence{}, at("00:00"), false},
|
|
}
|
|
for _, c := range cases {
|
|
if got := dueAt(c.c, c.t); got != c.want {
|
|
t.Errorf("%s: dueAt = %v, want %v", c.name, got, c.want)
|
|
}
|
|
}
|
|
}
|
|
|
|
// TestFrameworkScheduleSmoke covers the framework's own daily purge entry:
|
|
// first in the entry list when the app has config, removed by an empty
|
|
// purge_at, absent for a config-less app, and a boot error when malformed.
|
|
func TestFrameworkScheduleSmoke(t *testing.T) {
|
|
plugins := schedulePlugins(pact.ScheduledCommand{Command: "acme:tick", Cadence: pact.Daily()})
|
|
withConfig := func(t *testing.T, extra map[string]any) *backpack.App {
|
|
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)
|
|
}
|
|
}
|
|
return backpack.New(cfg)
|
|
}
|
|
|
|
entries, err := scheduleEntries(withConfig(t, nil), plugins)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(entries) != 2 || entries[0].id != "summercms.lagoon[0]:deferred:purge" || entries[0].plugin != "summercms.lagoon" {
|
|
t.Fatalf("entries = %+v", entries)
|
|
}
|
|
if entries[0].cmd.Cadence != pact.DailyAt(3, 0) || entries[1].id != "acme.test[0]:acme:tick" {
|
|
t.Fatalf("framework entry %+v, next %s", entries[0].cmd, entries[1].id)
|
|
}
|
|
_, table, err := periodicJobs(withConfig(t, map[string]any{"database.deferred_bindings.purge_at": "04:30"}), plugins)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if fw, ok := table["summercms.lagoon[0]:deferred:purge"]; !ok || fw.Cadence != pact.DailyAt(4, 30) {
|
|
t.Fatalf("compiled table = %+v", table)
|
|
}
|
|
|
|
entries, err = scheduleEntries(withConfig(t, map[string]any{"database.deferred_bindings.purge_at": ""}), plugins)
|
|
if err != nil || len(entries) != 1 || entries[0].id != "acme.test[0]:acme:tick" {
|
|
t.Fatalf("disabled: %+v (%v)", entries, err)
|
|
}
|
|
|
|
entries, err = scheduleEntries(backpack.New(nil), plugins)
|
|
if err != nil || len(entries) != 1 || entries[0].plugin == "summercms.lagoon" {
|
|
t.Fatalf("nil config: %+v (%v)", entries, err)
|
|
}
|
|
|
|
_, err = scheduleEntries(withConfig(t, map[string]any{"database.deferred_bindings.purge_at": "3am"}), plugins)
|
|
if err == nil || !strings.Contains(err.Error(), "database.deferred_bindings.purge_at") {
|
|
t.Fatalf("malformed purge_at = %v", err)
|
|
}
|
|
}
|