Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions apps/core/entities/job/job.entity.yml
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,10 @@ fields:
default: default
index: true

- name: cron
label: Cron
type: text

- name: timeout
label: Timeout
type: text
Expand Down
3 changes: 3 additions & 0 deletions docs/jobs.md
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ The Job key comes from the bundle folder name. There is no `name` field in `job.
label: Send Welcome Email
description: Sends the first welcome email to a new contact.
queue: default
cron: "0 9 * * MON"
timeout: 30s
retry:
attempts: 3
Expand All @@ -73,6 +74,7 @@ retry:
| `label` | yes | Human-facing Job label. |
| `description` | no | Human-facing explanation of the Job. |
| `queue` | no | Registered queue name. Missing `queue` uses `default`. |
| `cron` | no | Standard 5-field UTC cron expression. When set, the Job runs on that schedule. |
| `timeout` | yes | Positive Go duration string used for the handler deadline and worker lease. |
| `retry` | no | Retry settings. Missing `retry` means one attempt only. |
| `retry.attempts` | yes when `retry` exists | Total attempts, including the first try. Must be at least `2`. |
Expand All @@ -83,6 +85,7 @@ Rules:

- Job keys and queue names use kebab-case.
- Durations use Go duration syntax, such as `30s`, `5m`, or `1h30m`.
- `cron` uses standard 5-field syntax in UTC, such as `0 9 * * MON`. It is optional; omit it for manually enqueued Jobs.
- `retry.max-delay` must be greater than or equal to `retry.initial-delay`.
- Retry strategy is exponential and is not configurable in `job.yml`.
- Payloads are JSON. dygo stores them as JSON and handler code validates the business shape.
Expand Down
10 changes: 10 additions & 0 deletions docs/schedule.md
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,16 @@ dygo job execution list

`dygo worker --once` checks due Schedules, claims one available Job Execution batch, persists the result, and exits.

## Job Cron Shortcut

For the common case where one Job has one UTC schedule, put the cron expression directly in that Job's `job.yml`:

```yaml
cron: "0 9 * * MON"
```

Metadata sync creates a file-backed Schedule with the key `job-<job-key>` in the same App. Do not use that key in `_schedules.yml` while the Job has `cron`; duplicate keys fail metadata validation. If you remove `cron`, the next sync retires the generated Schedule. The worker then uses the same durable Schedule and Job Execution path described above. Use `_schedules.yml` when a Job needs a timezone other than UTC or more than one schedule.

## Coming Soon

- Studio UI for creating and editing Schedules.
Expand Down
11 changes: 7 additions & 4 deletions internal/db/metadata_records.go
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ type jobRecord struct {
Label string
Description string
Queue string
Cron string
Timeout string
Retry []byte
Enabled bool
Expand Down Expand Up @@ -354,20 +355,21 @@ func persistJobRecord(ctx context.Context, tx pgx.Tx, appID int64, job jobRecord
}
var id int64
err := tx.QueryRow(ctx, `
INSERT INTO "job" (name, app_id, key, source, label, description, queue, timeout, retry, enabled, retired)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
INSERT INTO "job" (name, app_id, key, source, label, description, queue, cron, timeout, retry, enabled, retired)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12)
ON CONFLICT (app_id, key) DO UPDATE
SET name = EXCLUDED.name,
source = EXCLUDED.source,
label = EXCLUDED.label,
description = EXCLUDED.description,
queue = EXCLUDED.queue,
cron = EXCLUDED.cron,
timeout = EXCLUDED.timeout,
retry = EXCLUDED.retry,
retired = false,
updated_at = now()
WHERE "job"."source" = $12
RETURNING id`, job.Name, appID, job.Key, source, job.Label, nullIfEmpty(job.Description), job.Queue, job.Timeout, job.Retry, job.Enabled, job.Retired, jobs.JobSourceFile).Scan(&id)
WHERE "job"."source" = $13
RETURNING id`, job.Name, appID, job.Key, source, job.Label, nullIfEmpty(job.Description), job.Queue, nullIfEmpty(job.Cron), job.Timeout, job.Retry, job.Enabled, job.Retired, jobs.JobSourceFile).Scan(&id)
if err != nil && err != pgx.ErrNoRows {
return 0, fmt.Errorf("persist job metadata %s/%s: %w", job.AppName, job.Key, err)
}
Expand Down Expand Up @@ -711,6 +713,7 @@ func buildMetadataRecords(metadata metadataCatalog) (metadataRecordSet, error) {
Label: loaded.Job.Label,
Description: loaded.Job.Description,
Queue: loaded.Job.EffectiveQueue(),
Cron: loaded.Job.Cron,
Timeout: loaded.Job.Timeout,
Retry: retryJSON,
Enabled: true,
Expand Down
3 changes: 2 additions & 1 deletion internal/db/metadata_records_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -213,6 +213,7 @@ func TestBuildMetadataRecordsStoresJobMetadata(t *testing.T) {
Label: "Send Welcome Email",
Description: "Sends a welcome email.",
Queue: "email",
Cron: "0 9 * * MON",
Timeout: "30s",
Retry: &jobs.Retry{Attempts: 3},
},
Expand All @@ -226,7 +227,7 @@ func TestBuildMetadataRecordsStoresJobMetadata(t *testing.T) {
t.Fatalf("job records count = %d, want 1", len(records.Jobs))
}
job := records.Jobs[0]
if job.Name != "sales.send-welcome-email" || job.Key != "send-welcome-email" || job.Source != jobs.JobSourceFile || job.Label != "Send Welcome Email" || job.Queue != "email" || job.Timeout != "30s" || !job.Enabled || job.Retired {
if job.Name != "sales.send-welcome-email" || job.Key != "send-welcome-email" || job.Source != jobs.JobSourceFile || job.Label != "Send Welcome Email" || job.Queue != "email" || job.Cron != "0 9 * * MON" || job.Timeout != "30s" || !job.Enabled || job.Retired {
t.Fatalf("job record = %+v, want synced sales job metadata", job)
}
for _, want := range []string{`"attempts":3`, `"initial-delay":"10s"`, `"max-delay":"5m"`} {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,10 @@ fields:
default: default
index: true

- name: cron
label: Cron
type: text

- name: timeout
label: Timeout
type: text
Expand Down
11 changes: 11 additions & 0 deletions internal/jobs/jobs.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import (
"github.com/hapyco/dygo/internal/queues"
"github.com/hapyco/dygo/internal/shape"
"github.com/hapyco/dygo/internal/yamlmeta"
"github.com/robfig/cron/v3"
"gopkg.in/yaml.v3"
)

Expand All @@ -24,6 +25,8 @@ const (
defaultRetryMaxDelay = "5m"
)

var jobCronParser = cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)

const (
// JobSourceFile marks Jobs synced from apps/<app>/jobs/<job>/job.yml.
JobSourceFile = "file"
Expand Down Expand Up @@ -55,6 +58,7 @@ type Job struct {
Label string `yaml:"label"`
Description string `yaml:"description,omitempty"`
Queue string `yaml:"queue,omitempty"`
Cron string `yaml:"cron,omitempty"`
Timeout string `yaml:"timeout"`
Retry *Retry `yaml:"retry,omitempty"`
}
Expand Down Expand Up @@ -227,6 +231,13 @@ func (j Job) Validate() error {
if queue := strings.TrimSpace(j.Queue); queue != "" && !fieldtype.IsName(queue) {
problems = append(problems, fmt.Sprintf("queue %q must be kebab-case", j.Queue))
}
if cronExpr := strings.TrimSpace(j.Cron); cronExpr != "" {
if len(strings.Fields(cronExpr)) != 5 {
problems = append(problems, "cron must contain exactly 5 fields in UTC, without CRON_TZ or TZ")
} else if _, err := jobCronParser.Parse(cronExpr); err != nil {
problems = append(problems, fmt.Sprintf("cron %q is invalid: %v", j.Cron, err))
}
}
if _, err := positiveDuration(j.Timeout, "timeout"); err != nil {
problems = append(problems, err.Error())
}
Expand Down
25 changes: 25 additions & 0 deletions internal/jobs/jobs_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -94,3 +94,28 @@ func writeJobTestFile(t *testing.T, path string, body string) {
t.Fatalf("WriteFile(%s) error = %v", path, err)
}
}

func TestDecodeJobCron(t *testing.T) {
for _, test := range []struct {
cron string
valid bool
}{
{"", true},
{"0 9 * * MON", true},
{"CRON_TZ=UTC 0 9 * * MON", false},
{"TZ=Asia/Karachi 0 9 * * MON", false},
{"* * * * * *", false},
{"@every 1h", false},
{"0 25 * * *", false},
} {
t.Run(test.cron, func(t *testing.T) {
job, err := Decode([]byte("label: Report\ntimeout: 30s\ncron: \"" + test.cron + "\"\n"))
if (err == nil) != test.valid {
t.Fatalf("Decode() error = %v, want valid=%v", err, test.valid)
}
if err == nil && job.Cron != test.cron {
t.Fatalf("Cron = %q, want %q", job.Cron, test.cron)
}
})
}
}
18 changes: 18 additions & 0 deletions internal/schedules/schedules.go
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,24 @@ func (c Catalog) Discover() ([]LoadedSchedule, error) {
}
schedules = append(schedules, discovered...)
}
for _, loaded := range c.jobs {
if strings.TrimSpace(loaded.Job.Cron) == "" {
continue
}
schedules = append(schedules, LoadedSchedule{
AppName: loaded.AppName,
AppDir: loaded.AppDir,
Path: loaded.Path,
Schedule: Schedule{
Name: "job-" + loaded.Job.Name,
Label: loaded.Job.Label,
Description: "Schedule for Job " + loaded.AppName + "/" + loaded.Job.Name,
Cron: loaded.Job.Cron,
Timezone: "UTC",
Job: loaded.AppName + "/" + loaded.Job.Name,
},
})
}
sortSchedules(schedules)
return schedules, nil
}
Expand Down
40 changes: 40 additions & 0 deletions internal/schedules/schedules_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
package schedules

import (
"github.com/hapyco/dygo/internal/app/manifest"
"os"
"path/filepath"
"strings"
"testing"
"time"
Expand Down Expand Up @@ -147,3 +150,40 @@ func TestCatalogSortsSchedules(t *testing.T) {
}
}
}

func TestCatalogValidatesJobCronSchedules(t *testing.T) {
appDir := t.TempDir()
if err := os.MkdirAll(filepath.Join(appDir, "jobs"), 0o755); err != nil {
t.Fatal(err)
}
apps := []manifest.LoadedApp{{Dir: appDir, Manifest: manifest.Manifest{Name: "sales"}}}
loadedJobs := []jobs.LoadedJob{{AppName: "sales", Path: "jobs/report/job.yml", Job: jobs.Job{
Name: "report", Label: "Report", Cron: "0 9 * * MON",
}}}
catalog := New(apps, loadedJobs)
loaded, err := catalog.Validate()
if err != nil || len(loaded) != 1 {
t.Fatalf("Validate() = %v, %v, want one Schedule", loaded, err)
}
schedule := loaded[0].Schedule
if schedule.Name != "job-report" || schedule.Timezone != "UTC" || schedule.Job != "sales/report" || schedule.Cron != "0 9 * * MON" || !schedule.EffectiveEnabled() {
t.Fatalf("generated Schedule = %+v", schedule)
}
if err := os.WriteFile(filepath.Join(appDir, "jobs", "_schedules.yml"), []byte(`schedules:
- name: job-report
label: Conflicting Report
cron: "0 10 * * MON"
timezone: UTC
job: sales/report
`), 0o644); err != nil {
t.Fatal(err)
}
if _, err := catalog.Validate(); err == nil || !strings.Contains(err.Error(), "duplicates Schedule identity") {
t.Fatalf("Validate() error = %v, want duplicate Schedule rejection", err)
}
loadedJobs[0].Job.Cron = ""
loaded, err = New(apps, loadedJobs).Validate()
if err != nil || len(loaded) != 1 || loaded[0].Schedule.Cron != "0 10 * * MON" {
t.Fatalf("Validate() after removing Job cron = %v, %v, want only explicit Schedule", loaded, err)
}
}
4 changes: 4 additions & 0 deletions schemas/job.schema.json
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,10 @@
"default": "default",
"description": "Registered queue name. Omitted Jobs use default."
},
"cron": {
"type": "string",
"description": "Optional standard five-field UTC cron expression."
},
"timeout": {
"$ref": "#/$defs/duration",
"description": "Required handler deadline and worker lease duration."
Expand Down
Loading