From d9f939a1eaf3e15402b7f94f57156c8b6e541771 Mon Sep 17 00:00:00 2001 From: Jakub Zych Date: Tue, 29 Sep 2026 19:47:51 +0200 Subject: [PATCH] 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 --- examples/hello/main.go | 3 + internal/build/build.go | 3 + internal/build/build_test.go | 20 +++ modules/bonfire/README.md | 17 +++ modules/bonfire/call.go | 72 ++++++++++ modules/bonfire/call_test.go | 55 ++++++++ modules/conga/README.md | 30 +++- modules/conga/conga.go | 5 + modules/conga/listen_test.go | 2 +- modules/conga/schedule.go | 104 ++++++++++++++ modules/conga/schedule_test.go | 250 +++++++++++++++++++++++++++++++++ modules/conga/scheduler.go | 216 ++++++++++++++++++++++++++++ modules/conga/worker.go | 19 +++ modules/pact/README.md | 30 +++- modules/pact/capabilities.go | 76 +++++++++- 15 files changed, 890 insertions(+), 12 deletions(-) create mode 100644 modules/bonfire/call.go create mode 100644 modules/bonfire/call_test.go create mode 100644 modules/conga/schedule.go create mode 100644 modules/conga/schedule_test.go create mode 100644 modules/conga/scheduler.go diff --git a/examples/hello/main.go b/examples/hello/main.go index 3eceaf7..44e23cd 100644 --- a/examples/hello/main.go +++ b/examples/hello/main.go @@ -45,6 +45,9 @@ func run(args []string, out io.Writer) error { commands = append(commands, hasCommands.Commands()...) } } + if err := app.Publish(bonfire.NewCatalog(commands)); err != nil { + return err + } root, err := bonfire.NewRoot("hello", commands, out) if err != nil { return err diff --git a/internal/build/build.go b/internal/build/build.go index ca57540..6550ca0 100644 --- a/internal/build/build.go +++ b/internal/build/build.go @@ -121,6 +121,9 @@ func generateMain(m Manifest) ([]byte, error) { b.WriteString("\t\t\tcommands = append(commands, hasCommands.Commands()...)\n") b.WriteString("\t\t}\n") b.WriteString("\t}\n") + b.WriteString("\tif err := app.Publish(bonfire.NewCatalog(commands)); err != nil {\n") + b.WriteString("\t\treturn err\n") + b.WriteString("\t}\n") fmt.Fprintf(&b, "\troot, err := bonfire.NewRoot(%s, commands, out)\n", strconv.Quote(m.Binary)) b.WriteString("\tif err != nil {\n") b.WriteString("\t\treturn err\n") diff --git a/internal/build/build_test.go b/internal/build/build_test.go index 1b7d36b..e1a1e9c 100644 --- a/internal/build/build_test.go +++ b/internal/build/build_test.go @@ -142,6 +142,26 @@ func TestGenerateMainRegistersCongaRuntimeCommands(t *testing.T) { } } +// TestGenerateMainPublishesCommandCatalog: the generated main publishes the +// final command list as a *bonfire.Catalog after collecting plugin commands +// and before building the root, so scheduled runs can call any command. +func TestGenerateMainPublishesCommandCatalog(t *testing.T) { + mainSrc, err := generateMain(Manifest{Module: "example.com/app", Binary: "hello"}) + if err != nil { + t.Fatal(err) + } + line := []byte("if err := app.Publish(bonfire.NewCatalog(commands)); err != nil {") + if n := bytes.Count(mainSrc, line); n != 1 { + t.Fatalf("catalog publish count = %d, want 1\n%s", n, mainSrc) + } + publish := bytes.Index(mainSrc, line) + plugins := bytes.Index(mainSrc, []byte("hasCommands.Commands()")) + root := bytes.Index(mainSrc, []byte("bonfire.NewRoot(")) + if plugins < 0 || root < 0 || !(plugins < publish && publish < root) { + t.Fatalf("catalog must be published after plugin commands and before NewRoot:\n%s", mainSrc) + } +} + func TestWriteIfChangedSkipsIdenticalBytes(t *testing.T) { dir := t.TempDir() path := filepath.Join(dir, "out.go") diff --git a/modules/bonfire/README.md b/modules/bonfire/README.md index 684e995..e2eb87c 100644 --- a/modules/bonfire/README.md +++ b/modules/bonfire/README.md @@ -18,6 +18,7 @@ bonfire is the console layer of SummerCMS. Plugins and framework modules describ - Widgets: box-drawn tables (`bonfire.Output.Table`), a spinner around a function (`bonfire.Output.Spinner`) and a progress bar (`bonfire.Output.Progress`); both fall back to plain lines when output is not a terminal. - Prompts: `bonfire.Output.Ask`, `bonfire.Output.Confirm`, `bonfire.Output.Choice` and `bonfire.Output.Secret`, which reads a hidden value on a terminal. Prompts return their defaults when input ends, and `bonfire.Output.Confirm` returns its default without asking when the session is not interactive. - Injectable streams (`bonfire.NewRootIO`, `bonfire.NewOutput`) so commands can be tested against buffers. +- In-process calls: `bonfire.Call` runs a named command with arguments against any writer (Laravel `Artisan::call`), and `bonfire.Catalog` holds an application binary's final command list so code outside the Cobra root, such as the conga scheduler, can call any registered command. ## Usage @@ -63,6 +64,18 @@ func main() { } ``` +Running a command in-process, with empty stdin so prompts take their defaults: + +```go +catalog := bonfire.NewCatalog(commands) +if catalog.Has("blog:import") { + err := catalog.Call(ctx, "blog:import", []string{"--dry-run", "https://example.com/feed"}, os.Stdout) + if errors.Is(err, bonfire.ErrUnknownCommand) { + // not registered in this binary + } +} +``` + ## API reference | Identifier | Description | @@ -77,6 +90,10 @@ func main() { | `bonfire.NewRootIO` | `bonfire.NewRoot` with injected stdin, stdout and stderr. | | `bonfire.NewOutput` | Builds a `bonfire.Output` over the given streams, applying the terminal and color policy. | | `bonfire.ErrCommandName` | Returned when a plugin command name is not in `namespace:verb` form. | +| `bonfire.Call` | Runs one command of a slice by exact name with arguments, writing output to a writer; stdin is empty. | +| `bonfire.ErrUnknownCommand` | Returned by `bonfire.Call` when no command has the name. | +| `bonfire.Catalog` | An immutable copy of a binary's command list; the generated app main publishes one on the app. | +| `bonfire.NewCatalog` | Builds a `bonfire.Catalog` from a command slice. | ## Configuration diff --git a/modules/bonfire/call.go b/modules/bonfire/call.go new file mode 100644 index 0000000..46e716b --- /dev/null +++ b/modules/bonfire/call.go @@ -0,0 +1,72 @@ +package bonfire + +import ( + "context" + "errors" + "fmt" + "io" + "strings" +) + +// ErrUnknownCommand is returned by Call when no command has the given name. +var ErrUnknownCommand = errors.New("bonfire: unknown command") + +// Call runs the command named name in-process with args, writing its output +// to out. Stdin is empty, so prompts take their defaults. A name that no +// command has is ErrUnknownCommand. +func Call(ctx context.Context, commands []Command, name string, args []string, out io.Writer) error { + cmd, ok := findCommand(commands, name) + if !ok { + return fmt.Errorf("%w: %q", ErrUnknownCommand, name) + } + if out == nil { + out = io.Discard + } + root, err := NewRootIO(name, []Command{cmd}, strings.NewReader(""), out, out) + if err != nil { + return err + } + root.SetArgs(append([]string{name}, args...)) + if ctx == nil { + ctx = context.Background() + } + return root.ExecuteContext(ctx) +} + +// Catalog is the final command list of an application binary. The generated +// main publishes it on the app so code outside the cobra root, such as the +// scheduler, can run any registered command with Call. +type Catalog struct { + commands []Command +} + +// NewCatalog returns a Catalog holding a copy of commands. +func NewCatalog(commands []Command) *Catalog { + return &Catalog{commands: append([]Command(nil), commands...)} +} + +// Has reports whether the catalog holds a command named name. +func (c *Catalog) Has(name string) bool { + if c == nil { + return false + } + _, ok := findCommand(c.commands, name) + return ok +} + +// Call runs the named catalog command in-process; see the package Call. +func (c *Catalog) Call(ctx context.Context, name string, args []string, out io.Writer) error { + if c == nil { + return fmt.Errorf("%w: %q", ErrUnknownCommand, name) + } + return Call(ctx, c.commands, name, args, out) +} + +func findCommand(commands []Command, name string) (Command, bool) { + for _, c := range commands { + if c.Name == name { + return c, true + } + } + return Command{}, false +} diff --git a/modules/bonfire/call_test.go b/modules/bonfire/call_test.go new file mode 100644 index 0000000..690e9e1 --- /dev/null +++ b/modules/bonfire/call_test.go @@ -0,0 +1,55 @@ +package bonfire + +import ( + "bytes" + "context" + "errors" + "strings" + "testing" +) + +func TestCall(t *testing.T) { + var gotArgs []string + var gotFlag string + cmds := []Command{{ + Name: "acme:echo", + Flags: []Flag{{Name: "loud", Bare: true}}, + Args: []Arg{{Name: "word"}}, + Run: func(ctx context.Context, in Input, out Output) error { + gotArgs = in.Args() + gotFlag, _ = in.Flag("loud") + out.Printf("echo %s\n", strings.Join(in.Args(), " ")) + return nil + }, + }} + + var out bytes.Buffer + if err := Call(t.Context(), cmds, "acme:echo", []string{"--loud", "hello", "world"}, &out); err != nil { + t.Fatal(err) + } + if strings.Join(gotArgs, ",") != "hello,world" || gotFlag != "true" { + t.Fatalf("args = %v, loud = %q", gotArgs, gotFlag) + } + if out.String() != "echo hello world\n" { + t.Fatalf("out = %q", out.String()) + } + + err := Call(t.Context(), cmds, "acme:missing", nil, &out) + if !errors.Is(err, ErrUnknownCommand) || !strings.Contains(err.Error(), `"acme:missing"`) { + t.Fatalf("err = %v, want ErrUnknownCommand naming the command", err) + } + + cat := NewCatalog(cmds) + cmds[0].Name = "acme:renamed" + if !cat.Has("acme:echo") || cat.Has("acme:renamed") { + t.Fatal("catalog does not hold a copy of the commands") + } + out.Reset() + if err := cat.Call(t.Context(), "acme:echo", []string{"x"}, &out); err != nil || out.String() != "echo x\n" { + t.Fatalf("catalog call: out %q, err %v", out.String(), err) + } + var nilCat *Catalog + if nilCat.Has("acme:echo") || !errors.Is(nilCat.Call(t.Context(), "acme:echo", nil, &out), ErrUnknownCommand) { + t.Fatal("nil catalog must hold no commands") + } +} diff --git a/modules/conga/README.md b/modules/conga/README.md index ef0807c..fbb81a4 100644 --- a/modules/conga/README.md +++ b/modules/conga/README.md @@ -10,6 +10,8 @@ conga is the SummerCMS counterpart of the WinterCMS queue plus the apparatus job Every River client runs on the one `*sql.DB` pool that `lagoon` opens. A worker started by `conga.StartWorker` uses `riverdatabasesql.NewWithPgxListener`: all queries go through the shared pool, and only Postgres `LISTEN` goes through a dedicated pgx pool with a single connection, so a job committed by any process wakes the worker immediately instead of waiting for the poll interval. +The scheduler is the Go form of WinterCMS `registerSchedule`. Plugins declare recurring console commands through `pact.HasSchedule`; every worker turns them into River periodic jobs on wall-clock `conga.Daily` and `conga.Every` schedules in the `app.timezone` location. Each due run is a job on the `scheduled` queue that calls the command in-process through the app's `bonfire.Catalog`. + ## Features - River-free job declarations: `conga.Job` turns `func(ctx context.Context, args T) error` into a `pact.Job`; `conga.OnQueue`, `conga.MaxAttempts` and `conga.Timeout` set per-job defaults. A `pact.Job` not built by `conga.Job` is rejected with `conga.ErrNotCongaJob`. @@ -20,6 +22,8 @@ Every River client runs on the one `*sql.DB` pool that `lagoon` opens. A worker - Cancellation in two parts: `conga.Manager.CancelJob` is the outside cancel (sets `is_canceled` and `conga.StatusStopped`, then cancels the River job, so a queued job never starts and a running job's context is cancelled); `conga.Manager.StopJob` is what a job calls on its own row after `conga.Manager.CheckIfCanceled` reports true (status only, the WinterCMS `cancelJob`). - Outcome rules in the worker: an error on an attempt before the last leaves the row in progress so River can retry; the final failed attempt, or a recovered panic on it, sets `conga.StatusError` with the error text under the metadata key `error`; an error on a row that was stopped or cancelled cancels the River job instead. A job that returns nil without completing its row leaves it as it is. Skipped work is recorded as complete with `{"skipped": true}` metadata. - Workers: `conga.StartWorker` registers every plugin job and starts one River client; `conga.WorkerOptions.Queues` limits it to some queues, and an unknown queue is `conga.ErrUnknownQueue` listing the known ones. `conga.Worker.Stop` stops it gracefully and cancels running jobs when its context ends. `conga.StartServeWorker` is the variant the `serve` command uses: it starts nothing when `queue.work_in_serve` is false. An app without jobs still gets a worker that starts and idles. +- Scheduled commands: every worker carries one River periodic job per `pact.HasSchedule` entry, in plugin activation order then declaration order, with the id `[]:`. Only the elected leader enqueues. Each run is a `conga.ScheduledCommandArgs` job on `conga.QueueScheduled` with one attempt (an interrupted run is not retried; the next period runs normally) and unique by args within its cadence period, so a leader failover cannot double-enqueue a period. The worker runs a job only when its entry exists in the compiled schedule and its command and arguments match that entry exactly, so a forged `river_job` row cannot run an arbitrary command. An unregistered command, or an app with no published `bonfire.Catalog`, is logged at Warn (`schedule: command not registered; skipping`) and skipped without failing the worker or other entries. Command output is logged line by line at Info with a `command` attribute; failures are logged with the duration. +- Wall-clock schedules: `conga.Daily` fires at the next `hh:mm` in its location and `conga.Every` at the next multiple of its interval since local midnight, so a restart never delays a daily run by up to a day the way `river.PeriodicInterval(24h)` would. On a DST day `conga.Daily` keeps the wall-clock time. A schedule entry with an empty command, a zero cadence, a daily time out of range, or an interval under one second or not dividing 24h fails the worker start with an error naming the plugin id and entry index. - Commands: `conga.RuntimeCommands` adds `queue:work` and `queue:clear` to the application binary. ## Usage @@ -94,6 +98,19 @@ func importPosts(ctx context.Context, m *conga.Manager, files []string) error { } ``` +A plugin schedules one of its registered commands; the worker runs it: + +```go +var _ pact.HasSchedule = (*Plugin)(nil) + +func (p *Plugin) Schedule() []pact.ScheduledCommand { + return []pact.ScheduledCommand{ + {Command: "blog:prune-drafts", Cadence: pact.Daily()}, + {Command: "blog:sync-feed", Args: []string{"--quiet"}, Cadence: pact.Every(15 * time.Minute)}, + } +} +``` + A worker runs in the same process or in a separate one: ```go @@ -130,11 +147,15 @@ defer w.Stop(context.Background()) | `conga.Job` | Wraps a typed job function as a `pact.Job` that conga can run on River. | | `conga.JobOption` | Per-job option: `conga.OnQueue`, `conga.MaxAttempts`, `conga.Timeout`. | | `conga.JobID` | Returns the record row id of the job running in a context. | -| `conga.StartWorker` | Registers plugin jobs and starts a River worker client. | +| `conga.StartWorker` | Registers plugin jobs and starts a River worker client carrying the plugins' periodic schedule jobs. | | `conga.StartServeWorker` | The worker of the `serve` command; nil when `queue.work_in_serve` is false. | | `conga.WorkerOptions` | Selects the queues a worker runs. | | `conga.RuntimeCommands` | Returns the `queue:work` and `queue:clear` commands. | | `conga.Worker` | A running worker; `conga.Worker.Stop` stops it and `conga.Worker.Queues` lists its queues. | +| `conga.Daily` | `river.PeriodicSchedule` firing at `Hour:Minute` every day in `Loc` (UTC when nil). | +| `conga.Every` | `river.PeriodicSchedule` firing at every multiple of `Interval` since local midnight in `Loc` (UTC when nil). | +| `conga.ScheduledCommandArgs` | Args of one scheduled run: compiled `Entry` id, `Command` and `Args`; kind `summer.scheduled_command`. | +| `conga.QueueScheduled` | The `scheduled` queue that scheduled runs are inserted on. | | `conga.ErrNoDatabase` | The app has not published the shared database handles. | | `conga.ErrNotCongaJob` | A registered `pact.Job` was not built by `conga.Job`. | | `conga.ErrRegistrationClosed` | Registration was attempted while a worker runs. | @@ -149,7 +170,8 @@ Keys are read from the compass config (`config/queue.yaml`, or `SUMMER_QUEUE__.. | `queue.work_in_serve` | `true` | Whether the `serve` command runs the job worker in its own process. Set `false` when a separate `queue:work` process runs the jobs. | | `queue.max_attempts` | `3` | Attempts per job before its record becomes `conga.StatusError`, unless the job or dispatch sets its own. | | `queue.job_timeout` | `300` | Per-attempt deadline, in seconds or as a duration string such as `5m`, unless the job sets `conga.Timeout`. | -| `queue.queues.` | `default: 4` | Concurrent workers per queue. The worker runs these queues plus every queue a registered job names plus `default`. | +| `queue.queues.` | `default: 4`, `scheduled: 1` | Concurrent workers per queue. The worker runs these queues plus every queue a registered job names plus `default` and `scheduled`. | +| `app.timezone` | `UTC` | IANA location of `pact.Daily`, `pact.DailyAt` and `pact.Every` schedules. An unknown name fails the worker start. | | `database.dsn` | none (required) | Also opens the worker's single-connection `LISTEN` pool. With PgBouncer, that connection must use session pooling or go straight to Postgres; transaction pooling cannot hold a `LISTEN`. | ```yaml @@ -172,7 +194,7 @@ queues: ## Dependencies -- SummerCMS modules: [backpack](../backpack/README.md), [bouncer](../bouncer/README.md) (the dispatching principal), [lagoon](../lagoon/README.md) (the shared pool, `lagoon.JobsTable` and the migrations that create River's schema and `summer_jobs`), [pact](../pact/README.md), [party](../party/README.md). +- SummerCMS modules: [backpack](../backpack/README.md), [bonfire](../bonfire/README.md) (commands and the `bonfire.Catalog` scheduled runs call), [bouncer](../bouncer/README.md) (the dispatching principal), [lagoon](../lagoon/README.md) (the shared pool, `lagoon.JobsTable` and the migrations that create River's schema and `summer_jobs`), [pact](../pact/README.md), [party](../party/README.md). - Third-party: `github.com/riverqueue/river` v0.47.0 with its `riverdriver/riverdatabasesql` and `rivertype` modules. River is the Postgres job queue the project stack names: it gives transactional inserts on the shared `*sql.DB`, retries, stuck-job rescue, leader election for periodic work and `LISTEN`/`NOTIFY` wake-ups. `github.com/jackc/pgx/v5/pgxpool` opens the listener pool; `gorm.io/gorm` writes the record rows. ## Testing @@ -181,4 +203,4 @@ queues: go test ./modules/conga/... ``` -The tests start a `postgres:16-alpine` container through testcontainers-go and migrate a fresh database per test with `lagoon.Migrate`, so they need a running Docker daemon. `TestListenPickupLatency` sets a 30-second poll interval and checks that a job committed by a separate client is picked up in under one second, while a poll-only worker does not pick it up within two. `go test -short ./modules/conga/...` skips the database tests. +The tests start a `postgres:16-alpine` container through testcontainers-go and migrate a fresh database per test with `lagoon.Migrate`, so they need a running Docker daemon. `TestListenPickupLatency` sets a 30-second poll interval and checks that a job committed by a separate client is picked up in under one second, while a poll-only worker does not pick it up within two. `TestScheduleRunsCommand` checks that an `Every(1s)` entry runs its command through a periodic job and that an unregistered command is skipped with a Warn log. `go test -short ./modules/conga/...` skips the database tests. diff --git a/modules/conga/conga.go b/modules/conga/conga.go index c8fd321..5ae739b 100644 --- a/modules/conga/conga.go +++ b/modules/conga/conga.go @@ -45,6 +45,11 @@ type Manager struct { jobs map[string]congaJob inserter *river.Client[*sql.Tx] worker *river.Client[*sql.Tx] + + // scheduledJob is the built-in scheduled command job; schedule is the + // compiled entry table of the running worker, keyed by entry id. + scheduledJob pact.Job + schedule map[string]pact.ScheduledCommand } // From returns the app's Manager, publishing a new one on first use. diff --git a/modules/conga/listen_test.go b/modules/conga/listen_test.go index 6171b44..4a05451 100644 --- a/modules/conga/listen_test.go +++ b/modules/conga/listen_test.go @@ -779,7 +779,7 @@ func TestStartServeWorker(t *testing.T) { if err != nil || w == nil { t.Fatalf("default: worker %v, err %v", w, err) } - if got := w.Queues(); !reflect.DeepEqual(got, []string{"default"}) { + if got := w.Queues(); !reflect.DeepEqual(got, []string{"default", QueueScheduled}) { t.Fatalf("queues = %v", got) } if err := w.Stop(context.Background()); err != nil { diff --git a/modules/conga/schedule.go b/modules/conga/schedule.go new file mode 100644 index 0000000..4f0c544 --- /dev/null +++ b/modules/conga/schedule.go @@ -0,0 +1,104 @@ +package conga + +import ( + "fmt" + "strings" + "time" + + "git.golem15.com/golem15/summercms/modules/backpack" + "git.golem15.com/golem15/summercms/modules/pact" + "github.com/riverqueue/river" +) + +// Daily is a river.PeriodicSchedule that fires once a day at Hour:Minute:00 +// wall-clock time in Loc (UTC when nil). Unlike PeriodicInterval(24h) it +// does not drift when the scheduler restarts. +type Daily struct { + Hour, Minute int + Loc *time.Location +} + +// Next returns the first Hour:Minute:00 in Loc strictly after now. On a DST +// day a time inside the gap resolves the way time.Date normalizes it. +func (d Daily) Next(now time.Time) time.Time { + loc := orUTC(d.Loc) + n := now.In(loc) + t := time.Date(n.Year(), n.Month(), n.Day(), d.Hour, d.Minute, 0, 0, loc) + if !t.After(n) { + t = time.Date(n.Year(), n.Month(), n.Day()+1, d.Hour, d.Minute, 0, 0, loc) + } + return t +} + +// Every is a river.PeriodicSchedule that fires at every multiple of Interval +// since local midnight in Loc (UTC when nil). Interval must be positive. +type Every struct { + Interval time.Duration + Loc *time.Location +} + +// Next returns local midnight plus the smallest multiple of Interval +// strictly after now, and the next local midnight when that multiple falls +// on or after it. +func (e Every) Next(now time.Time) time.Time { + loc := orUTC(e.Loc) + n := now.In(loc) + midnight := time.Date(n.Year(), n.Month(), n.Day(), 0, 0, 0, 0, loc) + nextMidnight := time.Date(n.Year(), n.Month(), n.Day()+1, 0, 0, 0, 0, loc) + if e.Interval <= 0 { + return nextMidnight + } + k := n.Sub(midnight)/e.Interval + 1 + t := midnight.Add(k * e.Interval) + if !t.Before(nextMidnight) { + return nextMidnight + } + return t +} + +func orUTC(loc *time.Location) *time.Location { + if loc == nil { + return time.UTC + } + return loc +} + +// scheduleFor returns the River schedule of a cadence in loc and its period. +func scheduleFor(c pact.Cadence, loc *time.Location) (river.PeriodicSchedule, time.Duration, error) { + if c.IsZero() { + return nil, 0, fmt.Errorf("cadence is zero") + } + if h, m, ok := c.At(); ok { + if h < 0 || h > 23 || m < 0 || m > 59 { + return nil, 0, fmt.Errorf("daily time %02d:%02d is out of range", h, m) + } + return Daily{Hour: h, Minute: m, Loc: loc}, 24 * time.Hour, nil + } + d := c.Interval() + if d <= 0 { + return nil, 0, fmt.Errorf("interval %s is not positive", d) + } + if d < time.Second { + return nil, 0, fmt.Errorf("interval %s is shorter than one second", d) + } + if (24*time.Hour)%d != 0 { + return nil, 0, fmt.Errorf("interval %s does not divide 24h evenly", d) + } + return Every{Interval: d, Loc: loc}, d, nil +} + +// appLocation is the app.timezone location, UTC when unset. +func appLocation(app *backpack.App) (*time.Location, error) { + name := "" + if app != nil && app.Config != nil { + name = strings.TrimSpace(app.Config.String("app.timezone")) + } + if name == "" { + return time.UTC, nil + } + loc, err := time.LoadLocation(name) + if err != nil { + return nil, fmt.Errorf("conga: app.timezone %q: %w", name, err) + } + return loc, nil +} diff --git a/modules/conga/schedule_test.go b/modules/conga/schedule_test.go new file mode 100644 index 0000000..c3b8303 --- /dev/null +++ b/modules/conga/schedule_test.go @@ -0,0 +1,250 @@ +package conga + +import ( + "context" + "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/pact" + "git.golem15.com/golem15/summercms/modules/party" +) + +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") +} diff --git a/modules/conga/scheduler.go b/modules/conga/scheduler.go new file mode 100644 index 0000000..a16cb8d --- /dev/null +++ b/modules/conga/scheduler.go @@ -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 +// "[]:". +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) +} diff --git a/modules/conga/worker.go b/modules/conga/worker.go index 1aa2464..7481e51 100644 --- a/modules/conga/worker.go +++ b/modules/conga/worker.go @@ -55,6 +55,11 @@ func (w *Worker) Queues() []string { // from database.dsn, so jobs committed by any process are picked up without // waiting for the poll interval. The worker keeps running after ctx is // cancelled; stop it with Worker.Stop. +// +// Every worker also carries the River periodic jobs of plugins' +// pact.HasSchedule entries (only the elected leader enqueues them) and runs +// the built-in scheduled command job on QueueScheduled when that queue is +// selected. An invalid schedule entry or app.timezone fails the start. func StartWorker(ctx context.Context, app *backpack.App, plugins []party.Plugin, o WorkerOptions) (*Worker, error) { m, err := From(app) if err != nil { @@ -77,13 +82,23 @@ func StartWorker(ctx context.Context, app *backpack.App, plugins []party.Plugin, } s := settingsFromApp(app) log := loggerFromApp(app) + periodic, table, err := periodicJobs(app, plugins) + if err != nil { + return nil, err + } m.mu.Lock() defer m.mu.Unlock() if m.worker != nil { return nil, fmt.Errorf("conga: a worker is already running") } + if err := m.registerLocked(m.scheduledJobLocked()); err != nil { + return nil, fmt.Errorf("conga: register scheduled command job: %w", err) + } known := knownQueues(s, m.jobs) + if _, configured := s.queues[QueueScheduled]; !configured { + known[QueueScheduled] = scheduledMaxWorkers + } queues, err := selectQueues(known, o.Queues) if err != nil { return nil, err @@ -103,6 +118,9 @@ func StartWorker(ctx context.Context, app *backpack.App, plugins []party.Plugin, } cfg := baseConfig(s, log) cfg.Workers = workers + // Every worker carries the periodic jobs whatever its queue filter: + // only the elected leader enqueues them. + cfg.PeriodicJobs = periodic cfg.Queues = map[string]river.QueueConfig{} for name, n := range queues { cfg.Queues[name] = river.QueueConfig{MaxWorkers: n} @@ -135,6 +153,7 @@ func StartWorker(ctx context.Context, app *backpack.App, plugins []party.Plugin, return nil, fmt.Errorf("conga: start worker: %w", err) } m.worker = client + m.schedule = table return &Worker{m: m, client: client, listener: listener, queues: sortedQueueNames(queues)}, nil } diff --git a/modules/pact/README.md b/modules/pact/README.md index 4c1684a..a0cbbab 100644 --- a/modules/pact/README.md +++ b/modules/pact/README.md @@ -1,24 +1,25 @@ # pact -Capability interfaces that compiled plugins implement to contribute routes, config, migrations, middleware, commands, admin screens, translations, mail templates and jobs. +Capability interfaces that compiled plugins implement to contribute routes, config, migrations, middleware, commands, admin screens, translations, mail templates, jobs and scheduled commands. `import "git.golem15.com/golem15/summercms/modules/pact"` ## Overview -`pact` is the contract layer between plugins and the framework. It holds interfaces and plain data types, with no behaviour of its own. A plugin opts into a capability by implementing one of the `Has*` interfaces, and the framework package that owns the capability discovers it with a type assertion ([party](../party/README.md) for config, [surf](../surf/README.md) for routes and middleware, [lagoon](../lagoon/README.md) for migrations, [cabana](../cabana/README.md) for admin controllers). It replaces the `register*()` methods of a WinterCMS PluginBase (`registerPermissions`, `registerNavigation`, `registerSettings` and so on) with small, separately implementable interfaces. +`pact` is the contract layer between plugins and the framework. It holds interfaces and plain data types, with no behaviour of its own. A plugin opts into a capability by implementing one of the `Has*` interfaces, and the framework package that owns the capability discovers it with a type assertion ([party](../party/README.md) for config, [surf](../surf/README.md) for routes and middleware, [lagoon](../lagoon/README.md) for migrations, [cabana](../cabana/README.md) for admin controllers, [conga](../conga/README.md) for jobs and schedules). It replaces the `register*()` methods of a WinterCMS PluginBase (`registerPermissions`, `registerNavigation`, `registerSettings` and so on) with small, separately implementable interfaces. ## Features -- Plugin capability interfaces: `pact.HasRoutes`, `pact.HasConfig`, `pact.HasMigrations`, `pact.HasCommands`, `pact.HasModels`, `pact.HasJobs`, `pact.HasLang`, `pact.HasLangOverrides` and `pact.HasMailTemplates`. +- Plugin capability interfaces: `pact.HasRoutes`, `pact.HasConfig`, `pact.HasMigrations`, `pact.HasCommands`, `pact.HasModels`, `pact.HasJobs`, `pact.HasSchedule`, `pact.HasLang`, `pact.HasLangOverrides` and `pact.HasMailTemplates`. - HTTP contracts: the `pact.Router` group builder (implemented by surf), the `pact.Middleware` type, and named, parameterized (`name:param`) and house-envelope middleware through `pact.HasMiddleware`, `pact.HasMiddlewareFactories` and `pact.HasHouseMiddleware`. - Backend registration data: `pact.Permission`, `pact.NavigationItem` and `pact.SettingsItem`, exposed through `pact.HasPermissions`, `pact.HasNavigation` and `pact.HasSettings`. - Admin controller contracts: `pact.AdminController`, `pact.HasAdminControllers`, `pact.AdminAssets` (embedded Winter-shaped admin YAML), `pact.AdminPermissioned` and `pact.AdminRecordSource`. - Admin extension contracts, so a plugin extends the compiled admin SPA without a Node build: `pact.AdminClientAssets` (per-controller JS and CSS from the plugin's embedded `assets/` tree, Winter's `addJs`/`addCss`), `pact.HasAdminActions` with `pact.AdminAction`, `pact.AdminActionInput` and `pact.AdminActionResult` (named toolbar and widget actions whose routes, CSRF check, permissions and record scoping the framework owns), and `pact.AdminPartialData` (the curated view model a partial template renders). - Optional admin hooks a controller or model can implement: list and form query scoping (`pact.ListExtendQuery`, `pact.FormExtendQuery`), create, update and delete hooks (`pact.FormBeforeCreate`, `pact.FormAfterUpdate`, `pact.FormBeforeDelete` and their siblings), relation hooks (`pact.RelationExtendManageQuery`, `pact.RelationExtendOptionsQuery`, `pact.RelationBeforeLink`), filter scopes (`pact.FilterScope`, `pact.FilterOptions`) and dropdown options (`pact.DropdownOptionsProvider`). - A background job contract (`pact.Job`, `pact.JobArgs`) that does not depend on any queue library. +- A schedule contract: `pact.HasSchedule` returns `pact.ScheduledCommand` entries (a registered command name, its arguments and a `pact.Cadence` built with `pact.Daily`, `pact.DailyAt` or `pact.Every`), the Go form of WinterCMS `registerSchedule`. It does not depend on any queue library either. - `pact.OptionalMessage`, a service an optional plugin can publish so others integrate with it without importing its package. -- `pact.HasModels`, `pact.HasJobs` and `pact.OptionalMessage` are declared for plugins to implement, but no framework package consumes them yet. +- `pact.HasModels` and `pact.OptionalMessage` are declared for plugins to implement, but no framework package consumes them yet. ## Usage @@ -59,6 +60,19 @@ func listPosts(w http.ResponseWriter, r *http.Request) {} func showPost(w http.ResponseWriter, r *http.Request) {} ``` +A plugin schedules one of its registered commands by implementing `pact.HasSchedule`: + +```go +var _ pact.HasSchedule = (*Plugin)(nil) + +func (p *Plugin) Schedule() []pact.ScheduledCommand { + return []pact.ScheduledCommand{ + {Command: "blog:prune-drafts", Cadence: pact.Daily()}, + {Command: "blog:ping", Args: []string{"--quiet"}, Cadence: pact.Every(15 * time.Minute)}, + } +} +``` + ## API reference | Identifier | Description | @@ -75,6 +89,12 @@ func showPost(w http.ResponseWriter, r *http.Request) {} | `pact.HasModels` | Exposes GORM models. | | `pact.Job` | Background unit of work that receives `pact.JobArgs`. | | `pact.HasJobs` | Registers background jobs. | +| `pact.HasSchedule` | Declares recurring console commands (`Schedule() []pact.ScheduledCommand`); conga workers run them. | +| `pact.ScheduledCommand` | One schedule entry: `Command`, `Args` and `Cadence`. | +| `pact.Cadence` | Opaque run frequency with `IsZero`, `Interval` (24h for daily cadences) and `At` (hour and minute of a daily cadence). | +| `pact.Daily` | Cadence at 00:00 every day in the app timezone (Laravel `->daily()`). | +| `pact.DailyAt` | Cadence at a given hour and minute every day in the app timezone. | +| `pact.Every` | Cadence at every multiple of an interval since local midnight; the interval must be at least 1s and divide 24h. | | `pact.HasLang` | Ships translation YAML under `lang//.yaml`. | | `pact.HasLangOverrides` | Replaces or adds translations of any loaded namespace, including the framework's own. | | `pact.HasMailTemplates` | Ships mail templates and layout aliases. | @@ -97,7 +117,7 @@ func showPost(w http.ResponseWriter, r *http.Request) {} - SummerCMS modules: [bonfire](../bonfire/README.md) (the command type in `pact.HasCommands`). - Third-party: `github.com/go-gormigrate/gormigrate/v2`, `gorm.io/gorm`. -- Standard library: `context`, `io/fs`, `net/http`. +- Standard library: `context`, `io/fs`, `net/http`, `time`. ## Testing diff --git a/modules/pact/capabilities.go b/modules/pact/capabilities.go index 7620723..2e243fc 100644 --- a/modules/pact/capabilities.go +++ b/modules/pact/capabilities.go @@ -4,6 +4,7 @@ import ( "context" "io/fs" "net/http" + "time" "git.golem15.com/golem15/summercms/modules/bonfire" "github.com/go-gormigrate/gormigrate/v2" @@ -98,6 +99,77 @@ type HasJobs interface { Jobs() []Job } +type cadenceKind uint8 + +const ( + cadenceNone cadenceKind = iota + cadenceDaily + cadenceEvery +) + +// Cadence is how often a scheduled command runs. Build one with Daily, +// DailyAt or Every; the zero value is no cadence and is rejected by the +// scheduler. +type Cadence struct { + kind cadenceKind + hour int + minute int + every time.Duration +} + +// Daily runs once a day at 00:00 in the app timezone (Laravel ->daily()). +func Daily() Cadence { return DailyAt(0, 0) } + +// DailyAt runs once a day at hour:minute in the app timezone. +func DailyAt(hour, minute int) Cadence { + return Cadence{kind: cadenceDaily, hour: hour, minute: minute} +} + +// Every runs at every multiple of d since local midnight in the app +// timezone. The scheduler requires d to be at least one second and to divide +// 24h evenly. +func Every(d time.Duration) Cadence { + return Cadence{kind: cadenceEvery, every: d} +} + +// IsZero reports whether c is the zero Cadence. +func (c Cadence) IsZero() bool { return c.kind == cadenceNone } + +// Interval is the period of c: 24h for a daily cadence, d for Every(d) and +// zero for the zero Cadence. +func (c Cadence) Interval() time.Duration { + switch c.kind { + case cadenceDaily: + return 24 * time.Hour + case cadenceEvery: + return c.every + } + return 0 +} + +// At returns the wall-clock time of a daily cadence; ok is false for Every +// and for the zero Cadence. +func (c Cadence) At() (hour, minute int, ok bool) { + if c.kind != cadenceDaily { + return 0, 0, false + } + return c.hour, c.minute, true +} + +// ScheduledCommand is one recurring run of a registered console command. +// Command is the command name (namespace:verb) and Args its arguments. +type ScheduledCommand struct { + Command string + Args []string + Cadence Cadence +} + +// HasSchedule is implemented by plugins that run console commands on a +// schedule. Only these compiled entries are ever executed by the scheduler. +type HasSchedule interface { + Schedule() []ScheduledCommand +} + // AdminController is the compile-time admin controller contract. Phase 9 // grows the schema pipeline; ID, model name and YAML config directory are // enough for generated stubs to compile. @@ -369,10 +441,10 @@ type OptionalMessage interface { // packages exist: // // HasListeners -// HasSchedule // // The kernel type-asserts HasConfig (party, before Register), HasCommands -// (generated app main, after Boot), HasMigrations (lagoon migrate), and +// (generated app main, after Boot), HasMigrations (lagoon migrate), +// HasJobs and HasSchedule (conga workers), and // HasMiddleware/HasMiddlewareFactories/HasHouseMiddleware/HasRoutes // (surf assemble). surf.BucketProvider is type-asserted in Assemble/ // BuildRouter (not a pact interface: pact cannot import surf without a cycle).