Files
summercms/modules/conga/commands.go
Jakub Zych 2237a640d2 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
2026-09-29 22:34:21 +02:00

226 lines
6.9 KiB
Go

package conga
import (
"context"
"fmt"
"os"
"os/signal"
"strconv"
"strings"
"syscall"
"time"
"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, schedule:run and queue:clear
// commands for the application binary.
func RuntimeCommands(app *backpack.App, plugins []party.Plugin) []bonfire.Command {
return []bonfire.Command{
{
Name: "queue:work",
Description: "Run background job workers in the foreground",
Flags: []bonfire.Flag{{
Name: "queue",
Description: "Queue to work (repeatable; default all known queues)",
Repeatable: true,
}},
Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error {
return withDB(ctx, app, func() error {
return work(ctx, app, plugins, in.Flags("queue"), out)
})
},
},
{
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.",
Args: []bonfire.Arg{{
Name: "queue",
Description: "The name of the queue to clear (default \"default\")",
}},
Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error {
queue, _ := in.Argument("queue")
queue = strings.TrimSpace(queue)
if queue == "" {
queue = defaultQueue
}
return withDB(ctx, app, func() error {
return clearQueue(ctx, app, queue, out)
})
},
},
}
}
func work(ctx context.Context, app *backpack.App, plugins []party.Plugin, queues []string, out bonfire.Output) error {
ctx, stop := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM)
defer stop()
w, err := StartWorker(ctx, app, plugins, WorkerOptions{Queues: queues})
if err != nil {
return err
}
out.Info("worker started on queues: " + strings.Join(w.Queues(), ", "))
<-ctx.Done()
stopCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
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 {
m, err := From(app)
if err != nil {
return err
}
client, err := m.insertClient()
if err != nil {
return err
}
out.Info(fmt.Sprintf("Clearing queue %q", queue))
total := 0
for {
res, err := client.JobDeleteMany(ctx, river.NewJobDeleteManyParams().
Queues(queue).
States(rivertype.JobStateAvailable, rivertype.JobStateScheduled, rivertype.JobStateRetryable).
First(clearBatch))
if err != nil {
return fmt.Errorf("conga: clear queue %q: %w", queue, err)
}
if len(res.Jobs) == 0 {
break
}
total += len(res.Jobs)
}
out.Info(fmt.Sprintf("Cleared %d jobs", total))
return nil
}
// withDB opens and publishes the shared pool for a CLI command and closes it
// when fn returns, the way lagoon's own commands do.
func withDB(ctx context.Context, app *backpack.App, fn func() error) error {
sqlDB, gdb, err := lagoon.OpenFromApp(ctx, app)
if err != nil {
return err
}
defer sqlDB.Close()
if err := lagoon.Publish(app, sqlDB, gdb); err != nil {
return err
}
return fn()
}