P2-01: schedule schema + CRUD API
The `schedules` table was already laid down in migration 0001; this
slice adds the Go-side data model, store CRUD with atomic version
bumps, and REST endpoints.
* `store.Schedule` + `RetentionPolicy` + `ScheduleOptions` typed
views (the wire form on the agent side keeps retention/options
as raw JSON since the agent just forwards them to restic).
* Store CRUD: CreateSchedule / GetSchedule / ListSchedulesByHost /
UpdateSchedule / DeleteSchedule. Each mutation bumps
`host_schedule_version` atomically in the same tx via UPSERT on
`host_schedule_version`. SetHostAppliedScheduleVersion records
what the agent has confirmed via schedule.ack (P2-02 will use it).
* REST endpoints under /api/hosts/{id}/schedules + /{sid}:
GET (list, with the version envelope so callers can detect
drift), POST (create), PUT (update — kind is immutable), DELETE.
* Validation: cron expressions parse via robfig/cron/v3 (same
parser the agent will use, so anything that validates here will
fire there); kind ∈ {backup, forget, prune, check} (init/unlock
are operator-only one-shot kinds, not schedulable); backup
schedules require ≥1 path; hooks rejected on non-backup kinds
(spec §14.3).
* All mutations audit-logged.
* Tests: store-level CRUD + version-bump invariants; REST happy
path (create→list→update→delete with version progression); REST
validation table covers each rejection code.
newTestServerWithHub now sets BootstrapToken so the schedules
handler tests can use the existing login flow without a parallel
test-server constructor.
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,280 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
// CreateSchedule inserts a new schedule and bumps the host's
|
||||
// schedule_version atomically. Returns the inserted row's
|
||||
// CreatedAt / UpdatedAt timestamps written into s.
|
||||
func (st *Store) CreateSchedule(ctx context.Context, s *Schedule) error {
|
||||
if s.ID == "" || s.HostID == "" {
|
||||
return errors.New("store: schedule id and host_id required")
|
||||
}
|
||||
now := time.Now().UTC()
|
||||
s.CreatedAt = now
|
||||
s.UpdatedAt = now
|
||||
if s.Paths == nil {
|
||||
s.Paths = []string{}
|
||||
}
|
||||
if s.Excludes == nil {
|
||||
s.Excludes = []string{}
|
||||
}
|
||||
if s.Tags == nil {
|
||||
s.Tags = []string{}
|
||||
}
|
||||
pathsJSON, _ := json.Marshal(s.Paths)
|
||||
excludesJSON, _ := json.Marshal(s.Excludes)
|
||||
tagsJSON, _ := json.Marshal(s.Tags)
|
||||
retentionJSON, _ := json.Marshal(s.RetentionPolicy)
|
||||
optionsJSON, _ := json.Marshal(s.Options)
|
||||
|
||||
tx, err := st.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("store: begin tx: %w", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
|
||||
if _, err := tx.ExecContext(ctx,
|
||||
`INSERT INTO schedules (
|
||||
id, host_id, kind, cron_expr, paths, excludes, tags,
|
||||
retention_policy, options, pre_hook, post_hook, enabled,
|
||||
created_at, updated_at
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
||||
s.ID, s.HostID, s.Kind, s.CronExpr,
|
||||
string(pathsJSON), string(excludesJSON), string(tagsJSON),
|
||||
string(retentionJSON), string(optionsJSON),
|
||||
s.PreHook, s.PostHook, boolToInt(s.Enabled),
|
||||
now.Format(time.RFC3339Nano), now.Format(time.RFC3339Nano),
|
||||
); err != nil {
|
||||
return fmt.Errorf("store: create schedule: %w", err)
|
||||
}
|
||||
if err := bumpHostScheduleVersionTx(ctx, tx, s.HostID); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
// UpdateSchedule replaces every editable field on an existing row
|
||||
// and bumps host_schedule_version. ID and HostID must match an
|
||||
// existing row; kind is immutable (creating a new schedule is
|
||||
// cheaper than re-keying retention/hooks).
|
||||
func (st *Store) UpdateSchedule(ctx context.Context, s *Schedule) error {
|
||||
if s.ID == "" || s.HostID == "" {
|
||||
return errors.New("store: schedule id and host_id required")
|
||||
}
|
||||
if s.Paths == nil {
|
||||
s.Paths = []string{}
|
||||
}
|
||||
if s.Excludes == nil {
|
||||
s.Excludes = []string{}
|
||||
}
|
||||
if s.Tags == nil {
|
||||
s.Tags = []string{}
|
||||
}
|
||||
pathsJSON, _ := json.Marshal(s.Paths)
|
||||
excludesJSON, _ := json.Marshal(s.Excludes)
|
||||
tagsJSON, _ := json.Marshal(s.Tags)
|
||||
retentionJSON, _ := json.Marshal(s.RetentionPolicy)
|
||||
optionsJSON, _ := json.Marshal(s.Options)
|
||||
now := time.Now().UTC()
|
||||
|
||||
tx, err := st.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("store: begin tx: %w", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
|
||||
res, err := tx.ExecContext(ctx,
|
||||
`UPDATE schedules SET
|
||||
cron_expr = ?, paths = ?, excludes = ?, tags = ?,
|
||||
retention_policy = ?, options = ?,
|
||||
pre_hook = ?, post_hook = ?, enabled = ?,
|
||||
updated_at = ?
|
||||
WHERE id = ? AND host_id = ?`,
|
||||
s.CronExpr,
|
||||
string(pathsJSON), string(excludesJSON), string(tagsJSON),
|
||||
string(retentionJSON), string(optionsJSON),
|
||||
s.PreHook, s.PostHook, boolToInt(s.Enabled),
|
||||
now.Format(time.RFC3339Nano),
|
||||
s.ID, s.HostID,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("store: update schedule: %w", err)
|
||||
}
|
||||
n, _ := res.RowsAffected()
|
||||
if n == 0 {
|
||||
return ErrNotFound
|
||||
}
|
||||
s.UpdatedAt = now
|
||||
if err := bumpHostScheduleVersionTx(ctx, tx, s.HostID); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
// DeleteSchedule removes a schedule and bumps host_schedule_version.
|
||||
// Returns ErrNotFound if no row matched.
|
||||
func (st *Store) DeleteSchedule(ctx context.Context, hostID, scheduleID string) error {
|
||||
tx, err := st.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("store: begin tx: %w", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
|
||||
res, err := tx.ExecContext(ctx,
|
||||
`DELETE FROM schedules WHERE id = ? AND host_id = ?`,
|
||||
scheduleID, hostID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("store: delete schedule: %w", err)
|
||||
}
|
||||
n, _ := res.RowsAffected()
|
||||
if n == 0 {
|
||||
return ErrNotFound
|
||||
}
|
||||
if err := bumpHostScheduleVersionTx(ctx, tx, hostID); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
// GetSchedule returns one schedule by (host_id, id). Returns
|
||||
// ErrNotFound on miss.
|
||||
func (st *Store) GetSchedule(ctx context.Context, hostID, scheduleID string) (*Schedule, error) {
|
||||
row := st.db.QueryRowContext(ctx,
|
||||
`SELECT id, host_id, kind, cron_expr, paths, excludes, tags,
|
||||
retention_policy, options, pre_hook, post_hook, enabled,
|
||||
created_at, updated_at
|
||||
FROM schedules WHERE id = ? AND host_id = ?`,
|
||||
scheduleID, hostID)
|
||||
s, err := scanSchedule(row)
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return nil, ErrNotFound
|
||||
}
|
||||
return s, err
|
||||
}
|
||||
|
||||
// ListSchedulesByHost returns every schedule for a host, ordered
|
||||
// by created_at. Empty slice on miss (not an error).
|
||||
func (st *Store) ListSchedulesByHost(ctx context.Context, hostID string) ([]Schedule, error) {
|
||||
rows, err := st.db.QueryContext(ctx,
|
||||
`SELECT id, host_id, kind, cron_expr, paths, excludes, tags,
|
||||
retention_policy, options, pre_hook, post_hook, enabled,
|
||||
created_at, updated_at
|
||||
FROM schedules WHERE host_id = ? ORDER BY created_at`,
|
||||
hostID)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("store: list schedules: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
out := []Schedule{}
|
||||
for rows.Next() {
|
||||
s, err := scanScheduleRow(rows)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out = append(out, *s)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// GetHostScheduleVersion returns the current version for a host,
|
||||
// or 0 if no row exists yet.
|
||||
func (st *Store) GetHostScheduleVersion(ctx context.Context, hostID string) (int64, error) {
|
||||
var v int64
|
||||
err := st.db.QueryRowContext(ctx,
|
||||
`SELECT version FROM host_schedule_version WHERE host_id = ?`, hostID).Scan(&v)
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return 0, nil
|
||||
}
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("store: get schedule version: %w", err)
|
||||
}
|
||||
return v, nil
|
||||
}
|
||||
|
||||
// SetHostAppliedScheduleVersion records the version the agent has
|
||||
// confirmed via schedule.ack. Idempotent.
|
||||
func (st *Store) SetHostAppliedScheduleVersion(ctx context.Context, hostID string, version int64) error {
|
||||
_, err := st.db.ExecContext(ctx,
|
||||
`UPDATE hosts SET applied_schedule_version = ? WHERE id = ?`,
|
||||
version, hostID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("store: set applied schedule version: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// bumpHostScheduleVersionTx upserts host_schedule_version, +1 each
|
||||
// call. Caller owns the tx.
|
||||
func bumpHostScheduleVersionTx(ctx context.Context, tx *sql.Tx, hostID string) error {
|
||||
if _, err := tx.ExecContext(ctx,
|
||||
`INSERT INTO host_schedule_version (host_id, version)
|
||||
VALUES (?, 1)
|
||||
ON CONFLICT(host_id) DO UPDATE SET version = version + 1`,
|
||||
hostID); err != nil {
|
||||
return fmt.Errorf("store: bump schedule version: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ----- scan helpers --------------------------------------------------
|
||||
|
||||
func scanSchedule(row *sql.Row) (*Schedule, error) {
|
||||
return scanScheduleRow(row)
|
||||
}
|
||||
|
||||
type scheduleScanner interface {
|
||||
Scan(dest ...any) error
|
||||
}
|
||||
|
||||
func scanScheduleRow(s scheduleScanner) (*Schedule, error) {
|
||||
var (
|
||||
out Schedule
|
||||
paths, excludes, tags, retention, options string
|
||||
createdAt, updatedAt string
|
||||
enabled int
|
||||
)
|
||||
err := s.Scan(&out.ID, &out.HostID, &out.Kind, &out.CronExpr,
|
||||
&paths, &excludes, &tags, &retention, &options,
|
||||
&out.PreHook, &out.PostHook, &enabled,
|
||||
&createdAt, &updatedAt)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if paths != "" {
|
||||
_ = json.Unmarshal([]byte(paths), &out.Paths)
|
||||
}
|
||||
if excludes != "" {
|
||||
_ = json.Unmarshal([]byte(excludes), &out.Excludes)
|
||||
}
|
||||
if tags != "" {
|
||||
_ = json.Unmarshal([]byte(tags), &out.Tags)
|
||||
}
|
||||
if retention != "" {
|
||||
_ = json.Unmarshal([]byte(retention), &out.RetentionPolicy)
|
||||
}
|
||||
if options != "" {
|
||||
_ = json.Unmarshal([]byte(options), &out.Options)
|
||||
}
|
||||
out.Enabled = enabled != 0
|
||||
if t, err := time.Parse(time.RFC3339Nano, createdAt); err == nil {
|
||||
out.CreatedAt = t
|
||||
}
|
||||
if t, err := time.Parse(time.RFC3339Nano, updatedAt); err == nil {
|
||||
out.UpdatedAt = t
|
||||
}
|
||||
return &out, nil
|
||||
}
|
||||
|
||||
func boolToInt(b bool) int {
|
||||
if b {
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
@@ -0,0 +1,122 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// makeSchedHost is a minimal host row to hang schedule tests off.
|
||||
func makeSchedHost(t *testing.T, s *Store) string {
|
||||
t.Helper()
|
||||
const id = "01HSCHEDHOST0000000000000"
|
||||
if err := s.CreateHost(context.Background(), Host{
|
||||
ID: id, Name: "sched-host", OS: "linux", Arch: "amd64",
|
||||
AgentVersion: "dev", ResticVersion: "0.16.0", ProtocolVersion: 1,
|
||||
EnrolledAt: time.Now().UTC(),
|
||||
}, "tokenhash", ""); err != nil {
|
||||
t.Fatalf("create host: %v", err)
|
||||
}
|
||||
return id
|
||||
}
|
||||
|
||||
func TestSchedulesCRUDAndVersionBump(t *testing.T) {
|
||||
t.Parallel()
|
||||
s := openTestStore(t)
|
||||
ctx := context.Background()
|
||||
hostID := makeSchedHost(t, s)
|
||||
|
||||
// Initial version is 0 (no row).
|
||||
v, err := s.GetHostScheduleVersion(ctx, hostID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if v != 0 {
|
||||
t.Fatalf("initial version: got %d, want 0", v)
|
||||
}
|
||||
|
||||
keepLast := 7
|
||||
sched := Schedule{
|
||||
ID: "01SCHED000000000000000001", HostID: hostID,
|
||||
Kind: "backup", CronExpr: "0 3 * * *",
|
||||
Paths: []string{"/etc", "/home"},
|
||||
Tags: []string{"nightly"},
|
||||
RetentionPolicy: RetentionPolicy{KeepLast: &keepLast},
|
||||
Enabled: true,
|
||||
}
|
||||
if err := s.CreateSchedule(ctx, &sched); err != nil {
|
||||
t.Fatalf("create: %v", err)
|
||||
}
|
||||
if sched.CreatedAt.IsZero() || sched.UpdatedAt.IsZero() {
|
||||
t.Fatalf("Create should populate timestamps")
|
||||
}
|
||||
|
||||
v, _ = s.GetHostScheduleVersion(ctx, hostID)
|
||||
if v != 1 {
|
||||
t.Fatalf("version after create: got %d, want 1", v)
|
||||
}
|
||||
|
||||
// Round-trip read.
|
||||
got, err := s.GetSchedule(ctx, hostID, sched.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("get: %v", err)
|
||||
}
|
||||
if got.CronExpr != "0 3 * * *" || len(got.Paths) != 2 {
|
||||
t.Fatalf("round-trip lost data: %+v", got)
|
||||
}
|
||||
if got.RetentionPolicy.KeepLast == nil || *got.RetentionPolicy.KeepLast != 7 {
|
||||
t.Fatalf("retention round-trip: %+v", got.RetentionPolicy)
|
||||
}
|
||||
|
||||
// List sees it.
|
||||
list, err := s.ListSchedulesByHost(ctx, hostID)
|
||||
if err != nil || len(list) != 1 || list[0].ID != sched.ID {
|
||||
t.Fatalf("list: err=%v rows=%v", err, list)
|
||||
}
|
||||
|
||||
// Update bumps version.
|
||||
sched.CronExpr = "*/30 * * * *"
|
||||
sched.Enabled = false
|
||||
if err := s.UpdateSchedule(ctx, &sched); err != nil {
|
||||
t.Fatalf("update: %v", err)
|
||||
}
|
||||
v, _ = s.GetHostScheduleVersion(ctx, hostID)
|
||||
if v != 2 {
|
||||
t.Fatalf("version after update: got %d, want 2", v)
|
||||
}
|
||||
got, _ = s.GetSchedule(ctx, hostID, sched.ID)
|
||||
if got.CronExpr != "*/30 * * * *" || got.Enabled {
|
||||
t.Fatalf("update did not persist: %+v", got)
|
||||
}
|
||||
|
||||
// Delete bumps version, returns ErrNotFound on second try.
|
||||
if err := s.DeleteSchedule(ctx, hostID, sched.ID); err != nil {
|
||||
t.Fatalf("delete: %v", err)
|
||||
}
|
||||
v, _ = s.GetHostScheduleVersion(ctx, hostID)
|
||||
if v != 3 {
|
||||
t.Fatalf("version after delete: got %d, want 3", v)
|
||||
}
|
||||
if err := s.DeleteSchedule(ctx, hostID, sched.ID); !errors.Is(err, ErrNotFound) {
|
||||
t.Fatalf("delete after delete: want ErrNotFound, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSetHostAppliedScheduleVersion(t *testing.T) {
|
||||
t.Parallel()
|
||||
s := openTestStore(t)
|
||||
ctx := context.Background()
|
||||
hostID := makeSchedHost(t, s)
|
||||
|
||||
if err := s.SetHostAppliedScheduleVersion(ctx, hostID, 42); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
host, err := s.GetHost(ctx, hostID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if host.AppliedScheduleVersion != 42 {
|
||||
t.Fatalf("got %d, want 42", host.AppliedScheduleVersion)
|
||||
}
|
||||
}
|
||||
@@ -70,6 +70,47 @@ type Host struct {
|
||||
RepoInitialisedAt *time.Time
|
||||
}
|
||||
|
||||
// Schedule mirrors one row of the schedules table. JSON columns
|
||||
// (paths, excludes, tags, retention_policy, options) are decoded
|
||||
// into Go-native shapes for ergonomics; the wire form on the agent
|
||||
// side keeps retention_policy / options as raw JSON since the agent
|
||||
// just forwards them to restic.
|
||||
type Schedule struct {
|
||||
ID string
|
||||
HostID string
|
||||
Kind string
|
||||
CronExpr string
|
||||
Paths []string
|
||||
Excludes []string
|
||||
Tags []string
|
||||
RetentionPolicy RetentionPolicy
|
||||
Options ScheduleOptions
|
||||
PreHook string
|
||||
PostHook string
|
||||
Enabled bool
|
||||
CreatedAt time.Time
|
||||
UpdatedAt time.Time
|
||||
}
|
||||
|
||||
// RetentionPolicy is the typed view of `restic forget --keep-*`.
|
||||
// All fields nullable so empty == "no policy / keep everything".
|
||||
type RetentionPolicy struct {
|
||||
KeepLast *int `json:"keep_last,omitempty"`
|
||||
KeepHourly *int `json:"keep_hourly,omitempty"`
|
||||
KeepDaily *int `json:"keep_daily,omitempty"`
|
||||
KeepWeekly *int `json:"keep_weekly,omitempty"`
|
||||
KeepMonthly *int `json:"keep_monthly,omitempty"`
|
||||
KeepYearly *int `json:"keep_yearly,omitempty"`
|
||||
}
|
||||
|
||||
// ScheduleOptions covers per-schedule knobs that aren't core to the
|
||||
// command itself — currently bandwidth caps. Stored as JSON so
|
||||
// future fields don't churn the schema.
|
||||
type ScheduleOptions struct {
|
||||
LimitUploadKBps *int `json:"limit_upload_kbps,omitempty"`
|
||||
LimitDownloadKBps *int `json:"limit_download_kbps,omitempty"`
|
||||
}
|
||||
|
||||
// EnrollmentToken is the issuer's view of a one-time token. The
|
||||
// raw token is returned only at create time; the DB stores its hash.
|
||||
type EnrollmentToken struct {
|
||||
|
||||
Reference in New Issue
Block a user