feat(12.2-01): add deferred:purge, its daily framework schedule and relation child hooks
- 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
This commit is contained in:
@@ -23,6 +23,7 @@ The scheduler is the Go form of WinterCMS `registerSchedule`. Plugins declare re
|
||||
- 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 `<plugin id>[<index>]:<command>`. 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.
|
||||
- Framework schedule entries: before the plugin entries, every worker carries the framework's own entries from `lagoon.FrameworkSchedule`, listed under the plugin id `lagoon.FrameworkScheduleID`. Today that is the daily `deferred:purge` at `database.deferred_bindings.purge_at` (default `03:00`), with the id `summercms.lagoon[0]:deferred:purge`; it runs, appears in `schedule:run --once` and is protected by the compiled entry table exactly like a plugin entry. An empty `purge_at` removes it, and a malformed one fails the worker start.
|
||||
- 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`, `schedule:run` and `queue:clear` to the application binary. `schedule:run` runs a scheduler-only worker in its own process, and `schedule:run --once` runs the entries due in the current minute without River, for system cron.
|
||||
|
||||
@@ -172,6 +173,7 @@ Keys are read from the compass config (`config/queue.yaml`, or `SUMMER_QUEUE__..
|
||||
| `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.<name>` | `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.deferred_bindings.purge_at` | `03:00` | Daily time of the framework's `deferred:purge` entry (read by `lagoon.FrameworkSchedule`); an empty value disables the entry. |
|
||||
| `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
|
||||
|
||||
@@ -790,3 +790,56 @@ func TestScheduleDueAt(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
|
||||
"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"
|
||||
@@ -60,15 +61,35 @@ func (e scheduleEntry) construct() (river.JobArgs, *river.InsertOpts) {
|
||||
}
|
||||
}
|
||||
|
||||
// scheduleEntries compiles the pact.HasSchedule entries of plugins in
|
||||
// activation order, then declaration order. Entry ids are
|
||||
// "<plugin id>[<index>]:<command>".
|
||||
// scheduleEntries compiles the framework's own entries
|
||||
// (lagoon.FrameworkSchedule, listed under lagoon.FrameworkScheduleID) and
|
||||
// then 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
|
||||
framework, err := lagoon.FrameworkSchedule(app)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("conga: framework schedule: %w", err)
|
||||
}
|
||||
for i, sc := range framework {
|
||||
sched, period, err := scheduleFor(sc.Cadence, loc)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("conga: framework schedule entry %d (%s): %w", i, sc.Command, err)
|
||||
}
|
||||
sc.Args = slices.Clone(sc.Args)
|
||||
out = append(out, scheduleEntry{
|
||||
id: fmt.Sprintf("%s[%d]:%s", lagoon.FrameworkScheduleID, i, sc.Command),
|
||||
plugin: lagoon.FrameworkScheduleID,
|
||||
index: i,
|
||||
cmd: sc,
|
||||
schedule: sched,
|
||||
period: period,
|
||||
})
|
||||
}
|
||||
for _, p := range plugins {
|
||||
hs, ok := p.(pact.HasSchedule)
|
||||
if !ok {
|
||||
|
||||
Reference in New Issue
Block a user