From dc861fb4a8ca7be65018736eebbf1c3e69a7afe9 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 22 Sep 2026 19:48:05 +0200 Subject: [PATCH] Do not arm a scheduled step of a module held as found on an adopted node (hq ADR 0103) --- cmd/mesh-host/main.go | 6 +++- internal/apply/schedule.go | 14 ++++++++- internal/apply/schedule_test.go | 55 ++++++++++++++++++++++++++++++--- 3 files changed, 68 insertions(+), 7 deletions(-) diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index cbc7718..168ab4e 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -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 // scheduler is the one-shot CLI path, which exits rather than staying up to fire anything. 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)} diff --git a/internal/apply/schedule.go b/internal/apply/schedule.go index 3057559..ccede49 100644 --- a/internal/apply/schedule.go +++ b/internal/apply/schedule.go @@ -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 // 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. -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() defer s.mu.Unlock() @@ -96,6 +100,14 @@ func (s *Scheduler) Sync(d *declaration.Declaration) { if !ok || c.Schedule == "" { 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) if err != nil { // The declaration parser already refused a malformed cron before this runs, so a diff --git a/internal/apply/schedule_test.go b/internal/apply/schedule_test.go index fcca78c..715b0db 100644 --- a/internal/apply/schedule_test.go +++ b/internal/apply/schedule_test.go @@ -244,7 +244,7 @@ func TestAScheduledStepRunsWhenDueAndNotBefore(t *testing.T) { d := &declaration.Declaration{Version: 1, Resources: []declaration.Resource{ 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. 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.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{ scheduledContainer(t, "* * * * *"), - }}) + }}, nil) // 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. @@ -371,7 +371,7 @@ func TestASlowRunSkipsTheNextDueRunRatherThanStacking(t *testing.T) { s := NewScheduler(clock, run, log) s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{ scheduledContainer(t, "* * * * *"), // every minute - }}) + }}, nil) // 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)) @@ -414,7 +414,7 @@ func TestSyncForgetsAScheduleTheDeclarationNoLongerNames(t *testing.T) { s.Sync(&declaration.Declaration{Version: 1, Resources: []declaration.Resource{ scheduledContainer(t, "* * * * *"), - }}) + }}, nil) s.mu.Lock() have := len(s.jobs) s.mu.Unlock() @@ -427,10 +427,55 @@ func TestSyncForgetsAScheduleTheDeclarationNoLongerNames(t *testing.T) { parseTrusted(t, `{"declaration":1,"resources":[ {"id":"web","type":"container","name":"web","image":"`+pinned+`"} ]}`).Resources[0], - }}) + }}, nil) s.Advance(context.Background(), time.Date(2026, 9, 7, 12, 5, 5, 0, time.UTC)) s.Wait() if n := rec.fireCount(); n != 0 { 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") + } +}