- 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
226 lines
6.9 KiB
Go
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()
|
|
}
|