feat(11-02): run plugin schedules as River periodic jobs through bonfire.Call
- pact.HasSchedule with ScheduledCommand and Daily/DailyAt/Every cadences (no River import) - bonfire.Call, Catalog and ErrUnknownCommand for in-process command runs - conga Daily/Every wall-clock schedules in app.timezone, periodic jobs on every worker, scheduled queue (MaxAttempts 1, unique by args within the cadence period) - scheduled worker runs only entries matching the compiled table; unregistered commands are skipped with a Warn log - generated app main publishes bonfire.NewCatalog(commands); hello main regenerated
This commit is contained in:
216
modules/conga/scheduler.go
Normal file
216
modules/conga/scheduler.go
Normal file
@@ -0,0 +1,216 @@
|
||||
package conga
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"slices"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"git.golem15.com/golem15/summercms/modules/backpack"
|
||||
"git.golem15.com/golem15/summercms/modules/bonfire"
|
||||
"git.golem15.com/golem15/summercms/modules/pact"
|
||||
"git.golem15.com/golem15/summercms/modules/party"
|
||||
"github.com/riverqueue/river"
|
||||
)
|
||||
|
||||
// QueueScheduled is the queue scheduled command runs are inserted on.
|
||||
const QueueScheduled = "scheduled"
|
||||
|
||||
// scheduledMaxWorkers is the default concurrency of QueueScheduled.
|
||||
const scheduledMaxWorkers = 1
|
||||
|
||||
// ScheduledCommandArgs are the args of one scheduled command run. Entry is
|
||||
// the compiled schedule entry id; the worker runs the command only when
|
||||
// Command and Args match that entry exactly.
|
||||
type ScheduledCommandArgs struct {
|
||||
Entry string `json:"entry"`
|
||||
Command string `json:"command"`
|
||||
Args []string `json:"args"`
|
||||
}
|
||||
|
||||
// Kind is the River job kind of a scheduled command run.
|
||||
func (ScheduledCommandArgs) Kind() string { return "summer.scheduled_command" }
|
||||
|
||||
// scheduleEntry is one compiled pact.HasSchedule entry.
|
||||
type scheduleEntry struct {
|
||||
id string
|
||||
plugin string
|
||||
index int
|
||||
cmd pact.ScheduledCommand
|
||||
schedule river.PeriodicSchedule
|
||||
period time.Duration
|
||||
}
|
||||
|
||||
func (e scheduleEntry) args() ScheduledCommandArgs {
|
||||
return ScheduledCommandArgs{Entry: e.id, Command: e.cmd.Command, Args: slices.Clone(e.cmd.Args)}
|
||||
}
|
||||
|
||||
// construct is the periodic job constructor: one run on QueueScheduled, never
|
||||
// retried, unique per entry within one cadence period.
|
||||
func (e scheduleEntry) construct() (river.JobArgs, *river.InsertOpts) {
|
||||
return e.args(), &river.InsertOpts{
|
||||
Queue: QueueScheduled,
|
||||
MaxAttempts: 1,
|
||||
UniqueOpts: river.UniqueOpts{ByArgs: true, ByPeriod: e.period},
|
||||
}
|
||||
}
|
||||
|
||||
// scheduleEntries compiles the pact.HasSchedule entries of plugins in
|
||||
// activation order, then declaration order. Entry ids are
|
||||
// "<plugin id>[<index>]:<command>".
|
||||
func scheduleEntries(app *backpack.App, plugins []party.Plugin) ([]scheduleEntry, error) {
|
||||
loc, err := appLocation(app)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out []scheduleEntry
|
||||
for _, p := range plugins {
|
||||
hs, ok := p.(pact.HasSchedule)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
for i, sc := range hs.Schedule() {
|
||||
if strings.TrimSpace(sc.Command) == "" {
|
||||
return nil, fmt.Errorf("conga: plugin %s schedule entry %d: command is empty", p.ID(), i)
|
||||
}
|
||||
sched, period, err := scheduleFor(sc.Cadence, loc)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("conga: plugin %s schedule entry %d (%s): %w", p.ID(), i, sc.Command, err)
|
||||
}
|
||||
sc.Args = slices.Clone(sc.Args)
|
||||
out = append(out, scheduleEntry{
|
||||
id: fmt.Sprintf("%s[%d]:%s", p.ID(), i, sc.Command),
|
||||
plugin: p.ID(),
|
||||
index: i,
|
||||
cmd: sc,
|
||||
schedule: sched,
|
||||
period: period,
|
||||
})
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// periodicJobs returns the River periodic jobs of plugins' schedules and the
|
||||
// compiled entry table keyed by entry id.
|
||||
func periodicJobs(app *backpack.App, plugins []party.Plugin) ([]*river.PeriodicJob, map[string]pact.ScheduledCommand, error) {
|
||||
entries, err := scheduleEntries(app, plugins)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
jobs := make([]*river.PeriodicJob, 0, len(entries))
|
||||
table := make(map[string]pact.ScheduledCommand, len(entries))
|
||||
for _, e := range entries {
|
||||
jobs = append(jobs, river.NewPeriodicJob(e.schedule, e.construct, &river.PeriodicJobOpts{ID: e.id}))
|
||||
table[e.id] = e.cmd
|
||||
}
|
||||
return jobs, table, nil
|
||||
}
|
||||
|
||||
// scheduledJob is the built-in job that runs scheduled commands. It is
|
||||
// created once per Manager so re-registration is a no-op.
|
||||
func (m *Manager) scheduledJobLocked() pact.Job {
|
||||
if m.scheduledJob == nil {
|
||||
m.scheduledJob = Job[ScheduledCommandArgs](m.runScheduled, OnQueue(QueueScheduled), MaxAttempts(1))
|
||||
}
|
||||
return m.scheduledJob
|
||||
}
|
||||
|
||||
// runScheduled runs a scheduled command job. Only an entry of the compiled
|
||||
// table whose command and args match the job exactly runs, so a forged
|
||||
// river_job row cannot execute an arbitrary command (T-11-09).
|
||||
func (m *Manager) runScheduled(ctx context.Context, a ScheduledCommandArgs) error {
|
||||
log := loggerFromApp(m.app)
|
||||
m.mu.Lock()
|
||||
entry, ok := m.schedule[a.Entry]
|
||||
m.mu.Unlock()
|
||||
if !ok || entry.Command != a.Command || !slices.Equal(entry.Args, a.Args) {
|
||||
log.Warn("schedule: job does not match a compiled schedule entry; skipping",
|
||||
"entry", a.Entry, "command", a.Command)
|
||||
return nil
|
||||
}
|
||||
w := newLogWriter(log, entry.Command)
|
||||
_, err := callScheduled(ctx, m.app, log, entry.Command, entry.Args, w)
|
||||
w.Flush()
|
||||
return err
|
||||
}
|
||||
|
||||
// callScheduled runs one scheduled command through the app's
|
||||
// *bonfire.Catalog. A missing catalog or an unregistered command is logged
|
||||
// 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)
|
||||
return true, nil
|
||||
}
|
||||
start := time.Now()
|
||||
err := cat.Call(ctx, command, args, out)
|
||||
d := time.Since(start)
|
||||
if err != nil {
|
||||
log.Error("schedule: command failed", "command", command, "duration", d, "error", err)
|
||||
return false, fmt.Errorf("conga: scheduled command %s: %w", command, err)
|
||||
}
|
||||
log.Info("schedule: command finished", "command", command, "duration", d)
|
||||
return false, nil
|
||||
}
|
||||
|
||||
// logWriter logs each output line of a scheduled command at Info.
|
||||
type logWriter struct {
|
||||
log *slog.Logger
|
||||
command string
|
||||
mu sync.Mutex
|
||||
buf bytes.Buffer
|
||||
}
|
||||
|
||||
func newLogWriter(log *slog.Logger, command string) *logWriter {
|
||||
return &logWriter{log: log, command: command}
|
||||
}
|
||||
|
||||
func (w *logWriter) Write(p []byte) (int, error) {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
w.buf.Write(p)
|
||||
for {
|
||||
line, err := w.buf.ReadString('\n')
|
||||
if err != nil {
|
||||
// Incomplete line: keep it for the next Write or Flush.
|
||||
w.buf.Reset()
|
||||
w.buf.WriteString(line)
|
||||
break
|
||||
}
|
||||
w.emit(line)
|
||||
}
|
||||
return len(p), nil
|
||||
}
|
||||
|
||||
// Flush logs a trailing line that did not end in a newline.
|
||||
func (w *logWriter) Flush() {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
if w.buf.Len() > 0 {
|
||||
w.emit(w.buf.String())
|
||||
w.buf.Reset()
|
||||
}
|
||||
}
|
||||
|
||||
func (w *logWriter) emit(line string) {
|
||||
line = strings.TrimRight(line, "\r\n")
|
||||
if line == "" {
|
||||
return
|
||||
}
|
||||
w.log.Info(line, "command", w.command)
|
||||
}
|
||||
Reference in New Issue
Block a user