Do not arm a scheduled step of a module held as found on an adopted node (hq ADR 0103)
This commit is contained in:
@@ -843,7 +843,11 @@ func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.D
|
|||||||
// declared is forgotten — and after a host restart the first apply rebuilds them all. A nil
|
// declared is forgotten — and after a host restart the first apply rebuilds them all. A nil
|
||||||
// scheduler is the one-shot CLI path, which exits rather than staying up to fire anything.
|
// scheduler is the one-shot CLI path, which exits rather than staying up to fire anything.
|
||||||
if sched != nil {
|
if sched != nil {
|
||||||
sched.Sync(declared)
|
held := map[string]bool{}
|
||||||
|
for _, h := range updated.Held {
|
||||||
|
held[h.ID] = true
|
||||||
|
}
|
||||||
|
sched.Sync(declared, held)
|
||||||
}
|
}
|
||||||
|
|
||||||
report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw)}
|
report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw)}
|
||||||
|
|||||||
@@ -86,7 +86,11 @@ func NewScheduler(clock Clock, run Runner, log func(string)) *Scheduler {
|
|||||||
// A job whose declaration is unchanged keeps its place in the cadence — its next due time and
|
// A job whose declaration is unchanged keeps its place in the cadence — its next due time and
|
||||||
// whether a run is in flight — so an ordinary reconcile every few minutes does not keep resetting
|
// whether a run is in flight — so an ordinary reconcile every few minutes does not keep resetting
|
||||||
// the clock out from under a schedule and prevent it ever firing.
|
// the clock out from under a schedule and prevent it ever firing.
|
||||||
func (s *Scheduler) Sync(d *declaration.Declaration) {
|
// held is the ids this node holds as found — what an adopted node keeps until its module is taken
|
||||||
|
// (novox/hq ADR 0100). A step of a module not yet taken is not armed: run on its cadence it would
|
||||||
|
// work on the predecessor's data, under the predecessor's service, which is the one thing an
|
||||||
|
// adopted node must not do. Nil on a converged node, where nothing is held.
|
||||||
|
func (s *Scheduler) Sync(d *declaration.Declaration, held map[string]bool) {
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
defer s.mu.Unlock()
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
@@ -96,6 +100,14 @@ func (s *Scheduler) Sync(d *declaration.Declaration) {
|
|||||||
if !ok || c.Schedule == "" {
|
if !ok || c.Schedule == "" {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
if module, untaken := d.Adoption.UntakenModuleOf(c.Identity()); untaken || held[c.Identity()] {
|
||||||
|
if module == "" {
|
||||||
|
module = "its module"
|
||||||
|
}
|
||||||
|
s.log(fmt.Sprintf("scheduled step %s: not armed while %s is held as found on this node",
|
||||||
|
c.Identity(), module))
|
||||||
|
continue
|
||||||
|
}
|
||||||
cron, err := declaration.ParseCron(c.Schedule)
|
cron, err := declaration.ParseCron(c.Schedule)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// The declaration parser already refused a malformed cron before this runs, so a
|
// The declaration parser already refused a malformed cron before this runs, so a
|
||||||
|
|||||||
@@ -244,7 +244,7 @@ func TestAScheduledStepRunsWhenDueAndNotBefore(t *testing.T) {
|
|||||||
d := &declaration.Declaration{Version: 1, Resources: []declaration.Resource{
|
d := &declaration.Declaration{Version: 1, Resources: []declaration.Resource{
|
||||||
scheduledContainer(t, "* * * * *"), // every minute; next due 12:01:00
|
scheduledContainer(t, "* * * * *"), // every minute; next due 12:01:00
|
||||||
}}
|
}}
|
||||||
s.Sync(d)
|
s.Sync(d, nil)
|
||||||
|
|
||||||
// Not yet due: 12:00:45 is before 12:01:00, so nothing runs.
|
// Not yet due: 12:00:45 is before 12:01:00, so nothing runs.
|
||||||
s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 0, 45, 0, time.UTC))
|
s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 0, 45, 0, time.UTC))
|
||||||
@@ -306,7 +306,7 @@ func TestAFailedRunIsRecordedAndDoesNotFailAnything(t *testing.T) {
|
|||||||
s := NewScheduler(clock, run, log)
|
s := NewScheduler(clock, run, log)
|
||||||
s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{
|
s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{
|
||||||
scheduledContainer(t, "* * * * *"),
|
scheduledContainer(t, "* * * * *"),
|
||||||
}})
|
}}, nil)
|
||||||
|
|
||||||
// Fire a run that exits non-zero. Advance returns nothing — there is no error to fail an apply,
|
// Fire a run that exits non-zero. Advance returns nothing — there is no error to fail an apply,
|
||||||
// because the run happens outside any apply and outside the store.
|
// because the run happens outside any apply and outside the store.
|
||||||
@@ -371,7 +371,7 @@ func TestASlowRunSkipsTheNextDueRunRatherThanStacking(t *testing.T) {
|
|||||||
s := NewScheduler(clock, run, log)
|
s := NewScheduler(clock, run, log)
|
||||||
s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{
|
s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{
|
||||||
scheduledContainer(t, "* * * * *"), // every minute
|
scheduledContainer(t, "* * * * *"), // every minute
|
||||||
}})
|
}}, nil)
|
||||||
|
|
||||||
// First occurrence: 12:01 is due — starts a run that blocks in the runner.
|
// First occurrence: 12:01 is due — starts a run that blocks in the runner.
|
||||||
s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 1, 5, 0, time.UTC))
|
s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 1, 5, 0, time.UTC))
|
||||||
@@ -414,7 +414,7 @@ func TestSyncForgetsAScheduleTheDeclarationNoLongerNames(t *testing.T) {
|
|||||||
|
|
||||||
s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{
|
s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{
|
||||||
scheduledContainer(t, "* * * * *"),
|
scheduledContainer(t, "* * * * *"),
|
||||||
}})
|
}}, nil)
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
have := len(s.jobs)
|
have := len(s.jobs)
|
||||||
s.mu.Unlock()
|
s.mu.Unlock()
|
||||||
@@ -427,10 +427,55 @@ func TestSyncForgetsAScheduleTheDeclarationNoLongerNames(t *testing.T) {
|
|||||||
parseTrusted(t, `{"declaration":1,"resources":[
|
parseTrusted(t, `{"declaration":1,"resources":[
|
||||||
{"id":"web","type":"container","name":"web","image":"`+pinned+`"}
|
{"id":"web","type":"container","name":"web","image":"`+pinned+`"}
|
||||||
]}`).Resources[0],
|
]}`).Resources[0],
|
||||||
}})
|
}}, nil)
|
||||||
s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 5, 5, 0, time.UTC))
|
s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 5, 5, 0, time.UTC))
|
||||||
s.Wait()
|
s.Wait()
|
||||||
if n := rec.fireCount(); n != 0 {
|
if n := rec.fireCount(); n != 0 {
|
||||||
t.Errorf("a schedule the declaration no longer names still fired (%d run(s))", n)
|
t.Errorf("a schedule the declaration no longer names still fired (%d run(s))", n)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Defends novox/hq ADR 0103: a scheduled step of a module not yet taken is not armed — run on its
|
||||||
|
// cadence it would work on the predecessor's data.
|
||||||
|
func TestAScheduledStepOfAnUntakenModuleIsNotArmed(t *testing.T) {
|
||||||
|
clock := &fixedClock{now: time.Date(2026, 9, 7, 12, 0, 30, 0, time.UTC)}
|
||||||
|
rec := &recordingRun{}
|
||||||
|
var said []string
|
||||||
|
s := NewScheduler(clock, rec.run, func(line string) { said = append(said, line) })
|
||||||
|
|
||||||
|
d := &declaration.Declaration{Version: 1,
|
||||||
|
Adoption: &declaration.Adoption{Untaken: map[string][]string{"backups": {"sync"}}},
|
||||||
|
Resources: []declaration.Resource{
|
||||||
|
scheduledContainer(t, "* * * * *"),
|
||||||
|
}}
|
||||||
|
s.Sync(d, nil)
|
||||||
|
s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 1, 5, 0, time.UTC))
|
||||||
|
s.Wait()
|
||||||
|
if rec.fireCount() != 0 {
|
||||||
|
t.Errorf("a held module's scheduled step ran %d time(s)", rec.fireCount())
|
||||||
|
}
|
||||||
|
if len(said) == 0 || !strings.Contains(said[0], "held as found") {
|
||||||
|
t.Errorf("nothing said why the step was not armed: %v", said)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Held by id — a step whose own container was found on the machine.
|
||||||
|
s2 := NewScheduler(clock, rec.run, nil)
|
||||||
|
s2.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{
|
||||||
|
scheduledContainer(t, "* * * * *"),
|
||||||
|
}}, map[string]bool{"sync": true})
|
||||||
|
s2.Advance(context.Background(), time.Date(2026, 9, 7, 12, 1, 5, 0, time.UTC))
|
||||||
|
s2.Wait()
|
||||||
|
if rec.fireCount() != 0 {
|
||||||
|
t.Errorf("a held scheduled step ran %d time(s)", rec.fireCount())
|
||||||
|
}
|
||||||
|
|
||||||
|
// Taken: the same step is armed and runs.
|
||||||
|
taken := &declaration.Declaration{Version: 1, Adoption: &declaration.Adoption{Taken: []string{"backups"}},
|
||||||
|
Resources: []declaration.Resource{scheduledContainer(t, "* * * * *")}}
|
||||||
|
s.Sync(taken, nil)
|
||||||
|
s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 2, 5, 0, time.UTC))
|
||||||
|
s.Wait()
|
||||||
|
if rec.fireCount() == 0 {
|
||||||
|
t.Error("a taken module's scheduled step never ran")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user