diff --git a/internal/apply/apply.go b/internal/apply/apply.go index b2ff24f..2f04871 100644 --- a/internal/apply/apply.go +++ b/internal/apply/apply.go @@ -445,7 +445,9 @@ func ApplyMindingWindows( // already the new one (novox/hq 04-ISSUES/103). That needs the file applied before the // container, which is the declared order; a container declared ahead of its file sees the // change one apply late, and never misses it. - in := inputs{declares: map[string]string{}, known: &known, windows: windows} + // And one patience with the artifact store for the whole apply (novox/hq issue 291): a fetch it + // does not answer during a maintenance window waits for the window instead of failing. + in := inputs{declares: map[string]string{}, known: &known, windows: windows, away: newStoreAway(windows, log)} for _, resource := range d.Resources { in.declares[resource.Identity()] = declaredDigest(resource) } @@ -917,9 +919,9 @@ func applyOne(ctx context.Context, sys system.System, r declaration.Resource, ru case *declaration.User: return applyUser(ctx, sys, res, run, previous) case *declaration.Archive: - return applyArchive(ctx, res, previous) + return applyArchive(ctx, res, previous, in.away) case *declaration.Process: - return applyProcess(ctx, res, run, changed, previous) + return applyProcess(ctx, res, run, changed, previous, in.away) case *declaration.Action: return applyAction(ctx, res, run) case *declaration.Network: @@ -1829,6 +1831,9 @@ type inputs struct { // windows is where a maintenance window says which containers it holds still (novox/hq issue // 224). Nil minds none. windows *Windows + // away is the apply's patience with an artifact store that does not answer (novox/hq issue + // 291). Nil fetches once. + away *storeAway } // fileDigest is what a file the container reads holds, by digest. diff --git a/internal/apply/archive.go b/internal/apply/archive.go index c82f025..87fcb1e 100644 --- a/internal/apply/archive.go +++ b/internal/apply/archive.go @@ -35,11 +35,11 @@ import ( // enough for a desktop theme and small enough to notice. const maxArchive = 512 << 20 -func applyArchive(ctx context.Context, r *declaration.Archive, previous store.Applied) (Outcome, error) { +func applyArchive(ctx context.Context, r *declaration.Archive, previous store.Applied, away *storeAway) (Outcome, error) { out := begin(r) out.Action = "unchanged" - body, err := fetch(ctx, r.Source) + body, err := fetchMinding(ctx, r.Source, away) if err != nil { return out, err } @@ -202,7 +202,7 @@ func fetch(ctx context.Context, source string) ([]byte, error) { } defer response.Body.Close() if response.StatusCode != http.StatusOK { - return nil, fmt.Errorf("%s answered %s", source, response.Status) + return nil, statusError{source: source, code: response.StatusCode, status: response.Status} } body, err := io.ReadAll(io.LimitReader(response.Body, maxArchive+1)) if err != nil { diff --git a/internal/apply/maintenance_window.go b/internal/apply/maintenance_window.go index 8275e7a..fefcebb 100644 --- a/internal/apply/maintenance_window.go +++ b/internal/apply/maintenance_window.go @@ -42,6 +42,7 @@ package apply // service the mesh can never bring back. import ( + "context" "encoding/json" "errors" "fmt" @@ -51,6 +52,8 @@ import ( "strings" "syscall" "time" + + "github.com/novox/mesh-host/internal/store" ) // WindowsName is the directory beside the node's state that holds the open windows, one file each. @@ -127,6 +130,56 @@ func (w *Windows) open(step string, holds []string, now time.Time) error { return writeAtomically(w.fileFor(step), raw, 0o600) } +// How long a step about to open a window waits for an apply in flight on this machine to end, and how +// often it looks (novox/hq issue 291). Past the bound it opens the window anyway, said: the applies +// then wait for the window instead (store_away.go), which is the slower order, never a failed one. +var ( + applyWaitAtMost = 15 * time.Minute + applyPoll = time.Second +) + +// openAfterApplies opens step's window only once no apply is in flight on this machine (novox/hq issue +// 291): an apply that started before the window would otherwise reach the store's archives with its +// server already held still. The apply lock is held only while the record is written — every apply +// that starts afterwards reads the window and waits for the store — so a push arriving mid-window +// still never queues behind the window itself (issue 224's reason for not taking the lock). +func (w *Windows) openAfterApplies(ctx context.Context, step string, holds []string, log func(string)) error { + if w == nil { + return nil + } + stateDir := filepath.Dir(w.dir) + deadline := time.Now().Add(applyWaitAtMost) + waited := false + for { + release, took, err := store.TryLockIn(stateDir) + if err != nil { + return err + } + if took { + defer release() + if waited { + log(fmt.Sprintf("scheduled step %s: the apply in flight ended; its window opens now", step)) + } + return w.open(step, holds, time.Now()) + } + if !waited { + waited = true + log(fmt.Sprintf("scheduled step %s: an apply is in flight on this machine — the window opens once "+ + "it ends, so nothing it fetches finds %s held still (novox/hq issue 291)", step, strings.Join(holds, ", "))) + } + if time.Now().After(deadline) { + log(fmt.Sprintf("scheduled step %s: the apply in flight did not end within %s; the window opens "+ + "beside it, and it waits for the window", step, applyWaitAtMost)) + return w.open(step, holds, time.Now()) + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(applyPoll): + } + } +} + // close removes the record of step's window. Called after its containers are started again; a file // already gone is the state wanted. func (w *Windows) close(step string) error { diff --git a/internal/apply/process.go b/internal/apply/process.go index 9ad9e47..e924de3 100644 --- a/internal/apply/process.go +++ b/internal/apply/process.go @@ -39,11 +39,11 @@ var daemonRoot = "/var/lib/mesh/daemons" var unitDir = "/etc/systemd/system" func applyProcess(ctx context.Context, r *declaration.Process, run Runner, - changed map[string]bool, previous store.Applied) (Outcome, error) { + changed map[string]bool, previous store.Applied, away *storeAway) (Outcome, error) { out := begin(r) out.Action = "unchanged" - body, err := fetch(ctx, r.Source) + body, err := fetchMinding(ctx, r.Source, away) if err != nil { return out, err } diff --git a/internal/apply/schedule.go b/internal/apply/schedule.go index 32ba578..2d8e969 100644 --- a/internal/apply/schedule.go +++ b/internal/apply/schedule.go @@ -261,7 +261,9 @@ func (s *Scheduler) fire(ctx context.Context, j *scheduledJob) { s.mu.Unlock() // The wall clock, not the scheduler's: the cadence is the scheduler's to decide, but how // long a window is believed is read by applies in other processes, against theirs. - if err := windows.open(j.id, j.hold, time.Now()); err != nil { + // And not while an apply is in flight here (novox/hq issue 291): one that started before the + // window would reach the store's archives with its server stopped under it. + if err := windows.openAfterApplies(ctx, j.id, j.hold, s.log); err != nil { s.log(fmt.Sprintf("scheduled step %s: cannot record the window it holds %v still for, "+ "so an apply could undo it — not run: %v", j.id, j.hold, err)) return diff --git a/internal/apply/store_away.go b/internal/apply/store_away.go new file mode 100644 index 0000000..89a78b7 --- /dev/null +++ b/internal/apply/store_away.go @@ -0,0 +1,159 @@ +package apply + +// A fetch from the artifact store while the store is away on purpose (novox/hq issue 291, ADR 0189). +// +// Every apply fetches each archive and process it declares, to know its bytes by digest. Each night +// the store's collector holds the store's server still for its maintenance window — seconds today, +// minutes on a larger store — and an apply that reached its archives in that window failed every one +// of them with "connection refused" and held the machine until its next pass. Harmless once, and +// wrong: a planned window must not fail an apply anywhere. +// +// **What the apply does about it, in the simplest form that is right on every machine:** +// +// - On the store's own machine the window is known exactly: it is the window record every apply +// there already reads (maintenance_window.go). A fetch the store does not answer while a window +// is open on this machine waits for the windows to close — polled, bounded by windowWaitAtMost — +// and is tried again. +// - On every other machine the store's window is not known, and is not published: a fetch the +// store does not answer is retried with backoff for at most storeAwayAtMost, **once per apply** — +// the budget is shared by every fetch of that apply, so a store that is really down costs one +// bounded wait, not one per archive. The bound is many times the window the collector holds +// (twelve seconds measured), and whatever still fails is retried by the engine's next pass +// anyway: the bound only decides whether one pass fails, never whether the machine converges. +// +// Only a store that does not answer is waited for: refused, reset, cut short, or a gateway saying +// the server behind it is away. An answer — a missing blob, the wrong bytes — fails at once, as before. + +import ( + "context" + "errors" + "fmt" + "io" + "net" + "net/http" + "strings" + "sync" + "syscall" + "time" +) + +// How long an apply waits for the windows open on its own machine to close before a fetch the store +// does not answer fails. Longer than any collection this mesh runs; far shorter than windowAtMost, +// which is how long a record nothing closed is believed. +var windowWaitAtMost = 20 * time.Minute + +// How long, in all, one apply retries a store that does not answer when no window on this machine +// explains it: on every machine but the store's, the store's own window. +var storeAwayAtMost = 3 * time.Minute + +// How often an open window is looked at again, and the first and longest backoff outside one. +// Variables so a test runs the 03:30 sequence in milliseconds. +var ( + windowPoll = 2 * time.Second + storeBackoff = time.Second + storeBackoffMax = 15 * time.Second +) + +// statusError is the store answering with a status other than 200. +type statusError struct { + source string + code int + status string +} + +func (e statusError) Error() string { return fmt.Sprintf("%s answered %s", e.source, e.status) } + +// storeAwayErr is whether a fetch failed because nothing answered for the store, rather than because +// the store answered no. +func storeAwayErr(err error) bool { + var s statusError + if errors.As(err, &s) { + return s.code == http.StatusBadGateway || s.code == http.StatusServiceUnavailable || + s.code == http.StatusGatewayTimeout + } + if errors.Is(err, syscall.ECONNREFUSED) || errors.Is(err, syscall.ECONNRESET) || + errors.Is(err, io.ErrUnexpectedEOF) || errors.Is(err, io.EOF) { + return true + } + var n net.Error + return errors.As(err, &n) && n.Timeout() +} + +// storeAway is one apply's patience with a store that does not answer. +type storeAway struct { + windows *Windows + log func(string) + + mu sync.Mutex + windowWaited time.Duration + left time.Duration + said map[string]bool +} + +// newStoreAway is the patience of one apply on a machine whose windows are recorded in windows (nil: +// none known, as on every machine but the store's). +func newStoreAway(windows *Windows, log func(string)) *storeAway { + if log == nil { + log = func(string) {} + } + return &storeAway{windows: windows, log: log, left: storeAwayAtMost, said: map[string]bool{}} +} + +// next is how long to wait before fetching again after the attempt-th failure, or zero and why not. +func (a *storeAway) next(attempt int) (time.Duration, string) { + a.mu.Lock() + defer a.mu.Unlock() + if open := a.windows.Open(time.Now()); len(open) > 0 { + if a.windowWaited < windowWaitAtMost { + if key := "window " + open[0].Step; !a.said[key] { + a.said[key] = true + a.log(fmt.Sprintf(" waiting the store does not answer while a maintenance window is open "+ + "on this machine (%s): fetched again once it closes, for at most %s (novox/hq issue 291)", + open[0], windowWaitAtMost)) + } + a.windowWaited += windowPoll + return windowPoll, "" + } + return 0, fmt.Sprintf("waited %s for the maintenance window open on this machine (%s), and it did not "+ + "close", a.windowWaited.Round(time.Second), open[0]) + } + if a.left <= 0 { + return 0, fmt.Sprintf("the store did not answer for the %s this apply waits for it; the next pass "+ + "fetches it again", storeAwayAtMost) + } + wait := storeBackoff + for i := 1; i < attempt && wait < storeBackoffMax; i++ { + wait *= 2 + } + wait = min(wait, storeBackoffMax, a.left) + a.left -= wait + if !a.said["away"] { + a.said["away"] = true + a.log(fmt.Sprintf(" waiting the store does not answer — its maintenance window, perhaps, on "+ + "its own machine: fetched again for at most %s in this apply (novox/hq issue 291)", storeAwayAtMost)) + } + return wait, "" +} + +// fetchMinding is fetch, patient with a store that is away on purpose. A nil away is not patient: a +// caller outside an apply. +func fetchMinding(ctx context.Context, source string, away *storeAway) ([]byte, error) { + for attempt := 1; ; attempt++ { + body, err := fetch(ctx, source) + if err == nil || away == nil || !storeAwayErr(err) { + return body, err + } + wait, why := away.next(attempt) + if wait <= 0 { + if why != "" && !strings.Contains(err.Error(), why) { + err = fmt.Errorf("%w — %s", err, why) + } + return nil, err + } + select { + case <-ctx.Done(): + return nil, err + case <-time.After(wait): + } + } +} diff --git a/internal/apply/store_away_test.go b/internal/apply/store_away_test.go new file mode 100644 index 0000000..497f141 --- /dev/null +++ b/internal/apply/store_away_test.go @@ -0,0 +1,307 @@ +package apply + +import ( + "context" + "net" + "net/http" + "path/filepath" + "strings" + "sync" + "testing" + "time" + + "github.com/novox/mesh-host/internal/declaration" + "github.com/novox/mesh-host/internal/store" +) + +// The 03:30 sequence (novox/hq issue 291): the store's collector holds the store's server still for its +// maintenance window, and the node-engine's own reconcile on that machine — and every other machine's — +// reached the store's archives in that window. Every bundle failed "connection refused" and the machine +// was held until the next pass. A planned window must fail no apply: the step waits for an apply in +// flight before it opens, an apply on the store's machine waits for the window to close, and an apply +// elsewhere waits a bounded while for a store that does not answer. + +// aPretendStore is the artifact store's address: answering, or refusing connections as a stopped +// server does — on the same address, as the real one comes back on. +type aPretendStore struct { + t *testing.T + addr string + body []byte + code int + + mu sync.Mutex + srv *http.Server + asked int +} + +func newPretendStore(t *testing.T, body []byte) *aPretendStore { + t.Helper() + l, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + s := &aPretendStore{t: t, addr: l.Addr().String(), body: body, code: http.StatusOK} + s.serve(l) + t.Cleanup(s.stop) + return s +} + +func (s *aPretendStore) serve(l net.Listener) { + srv := &http.Server{Handler: http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + s.mu.Lock() + s.asked++ + code := s.code + s.mu.Unlock() + w.WriteHeader(code) + _, _ = w.Write(s.body) + })} + s.mu.Lock() + s.srv = srv + s.mu.Unlock() + go func() { _ = srv.Serve(l) }() +} + +// stop is the server held still: connections refused from here on. +func (s *aPretendStore) stop() { + s.mu.Lock() + srv := s.srv + s.srv = nil + s.mu.Unlock() + if srv != nil { + _ = srv.Close() + } +} + +// start is the server running again, on its address. +func (s *aPretendStore) start() { + s.t.Helper() + l, err := net.Listen("tcp", s.addr) + if err != nil { + s.t.Errorf("the store's address could not be taken again: %v", err) + return + } + s.serve(l) +} + +func (s *aPretendStore) url(name string) string { + return "http://" + s.addr + "/v2/" + name + "/blobs/sha256:x" +} + +func (s *aPretendStore) fetches() int { s.mu.Lock(); defer s.mu.Unlock(); return s.asked } + +// patiently makes every wait of this file milliseconds, and puts the real ones back after. +func patiently(t *testing.T, away time.Duration) { + t.Helper() + was := []time.Duration{windowWaitAtMost, storeAwayAtMost, windowPoll, storeBackoff, storeBackoffMax, + applyWaitAtMost, applyPoll} + windowWaitAtMost, storeAwayAtMost, windowPoll, storeBackoff, storeBackoffMax = 10*time.Second, away, + 5*time.Millisecond, time.Millisecond, 10*time.Millisecond + applyWaitAtMost, applyPoll = 10*time.Second, 5*time.Millisecond + t.Cleanup(func() { + windowWaitAtMost, storeAwayAtMost, windowPoll, storeBackoff, storeBackoffMax = was[0], was[1], was[2], + was[3], was[4] + applyWaitAtMost, applyPoll = was[5], was[6] + }) +} + +// theStoreWithABundle is the store's machine: its server, its nightly collector, and a bundle this +// machine runs, fetched from the store. +func theStoreWithABundle(t *testing.T, source, digest, dir string) *declaration.Declaration { + t.Helper() + return parseTrusted(t, `{"declaration":1,"resources":[ + {"id":"store","type":"container","name":"mesh-registry","image":"`+pinned+`"}, + {"id":"collect","type":"container","name":"mesh-registry-collect","image":"`+pinned+`", + "schedule":"30 3 * * *","while-stopped":["store"]}, + {"id":"postgres.bundle-code","type":"archive","source":"`+source+`","digest":"`+digest+`","path":"`+dir+`/code"} + ]}`) +} + +type saying struct { + mu sync.Mutex + lines []string +} + +func (s *saying) say(line string) { s.mu.Lock(); s.lines = append(s.lines, line); s.mu.Unlock() } +func (s *saying) all() string { s.mu.Lock(); defer s.mu.Unlock(); return strings.Join(s.lines, "\n") } + +// On the store's own machine: an apply that reaches the store's archives while the collector holds the +// server still waits for the window, and applies once it closes — no failure, nothing held. +func TestAnApplyOnTheStoresMachineWaitsForItsWindow(t *testing.T) { + patiently(t, time.Second) + body, digest := anArchive(t, map[string]string{"bin/postgres-tools": "#!/bin/sh"}) + st := newPretendStore(t, body) + dir := t.TempDir() + d := theStoreWithABundle(t, st.url("postgres/code"), digest, dir) + m := newRegistryMachine(containerSpec(d.Resources[0].(*declaration.Container), inputs{})) + windows := WindowsIn(t.TempDir()) + + // 03:30: the collector holds mesh-registry still — the store refuses connections. + s := openAWindow(t, d, m, windows) + st.stop() + + // The reconcile reaches the bundle in the window. + said := &saying{} + type result struct { + report Report + err error + } + done := make(chan result, 1) + go func() { + report, _, err := ApplyMindingWindows(context.Background(), archHost(t), d, store.State{}, + store.OriginDeclared, m.run, said.say, nil, nil, windows) + done <- result{report, err} + }() + select { + case r := <-done: + t.Fatalf("the apply ended inside the window instead of waiting for it: %v", r.err) + case <-time.After(100 * time.Millisecond): + } + if !strings.Contains(said.all(), "maintenance window is open on this machine") { + t.Fatalf("the wait is not said: %s", said.all()) + } + + // 03:30:30: the collection ends; the server runs again, and the window closes. + st.start() + close(m.stepGo) + s.Wait() + var r result + select { + case r = <-done: + case <-time.After(5 * time.Second): + t.Fatal("the apply did not go on once the window closed") + } + if r.err != nil { + t.Fatalf("a planned window failed the apply: %v", r.err) + } + if o := outcomeOf(r.report, "postgres.bundle-code"); o.Action != "created" { + t.Fatalf("the bundle fetched after the window is %q: %+v", o.Action, o) + } +} + +// And the order on the store's machine: a collector due while an apply is in flight opens its window only +// once the apply ended — so an apply that started before it never finds the server stopped under it. +func TestAWindowOpensOnlyOnceTheApplyInFlightEnded(t *testing.T) { + patiently(t, time.Second) + d := aStoreWithACollector(t) + m := newRegistryMachine(containerSpec(d.Resources[0].(*declaration.Container), inputs{})) + stateDir := t.TempDir() + windows := WindowsIn(stateDir) + + // An apply holds the machine: the daemon's reconcile, at 03:30:00. + release, err := store.Lock(filepath.Join(stateDir, "state.json"), nil) + if err != nil { + t.Fatal(err) + } + said := &saying{} + clock := &fixedClock{now: time.Date(2026, 10, 2, 3, 29, 0, 0, time.UTC)} + s := NewScheduler(clock, m.run, said.say) + s.RecordWindowsIn(windows) + s.Sync(d, nil) + s.Advance(context.Background(), time.Date(2026, 10, 2, 3, 30, 5, 0, time.UTC)) + + time.Sleep(100 * time.Millisecond) + m.mu.Lock() + early := append([]string{}, m.calls...) + m.mu.Unlock() + if len(early) != 0 || len(windows.Open(time.Now())) != 0 { + t.Fatalf("the window opened under an apply in flight: %v", early) + } + if !strings.Contains(said.all(), "an apply is in flight") { + t.Fatalf("the wait is not said: %s", said.all()) + } + + // The apply ends: the collector holds the server still, runs, and lets it go. + release() + select { + case <-m.stepIn: + case <-time.After(5 * time.Second): + t.Fatal("the step never ran once the apply ended") + } + if len(windows.Open(time.Now())) != 1 { + t.Fatal("the window is not recorded once it opened") + } + close(m.stepGo) + s.Wait() + m.mu.Lock() + defer m.mu.Unlock() + var window []string + for _, c := range m.calls { + if !strings.HasPrefix(c, "rm ") { + window = append(window, c) + } + } + if strings.Join(window, ",") != "stop mesh-registry,step,start mesh-registry" { + t.Fatalf("the window ran as %v", m.calls) + } +} + +// On every other machine the store's window is not known: a store that does not answer is fetched again +// with backoff, and the bundle arrives once the store is back. +func TestAnApplyElsewhereWaitsABoundedWhileForTheStore(t *testing.T) { + patiently(t, 5*time.Second) + body, digest := anArchive(t, map[string]string{"bin/tools": "#!/bin/sh"}) + st := newPretendStore(t, body) + st.stop() + go func() { time.Sleep(60 * time.Millisecond); st.start() }() + + dir := t.TempDir() + d := declare(t, `{"id":"postgres.bundle-code","type":"archive","source":"`+st.url("postgres/code")+ + `","digest":"`+digest+`","path":"`+dir+`/code"}`) + said := &saying{} + report, _, err := ApplyMindingWindows(context.Background(), archHost(t), d, store.State{}, + store.OriginDeclared, noServices, said.say, nil, nil, nil) + if err != nil { + t.Fatalf("a store away for its window failed the apply elsewhere: %v", err) + } + if o := outcomeOf(report, "postgres.bundle-code"); o.Action != "created" { + t.Fatalf("the bundle is %q", o.Action) + } + if !strings.Contains(said.all(), "the store does not answer") { + t.Fatalf("the wait is not said: %s", said.all()) + } +} + +// aDigest is a digest no fetched bytes will match: these bundles are never meant to arrive. +const aDigest = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + +// A store really down costs one bounded wait per apply, not one per bundle; and an answer — a blob the +// store does not have — fails at once, as it always did. +func TestAStoreThatStaysAwayCostsOneBoundedWait(t *testing.T) { + patiently(t, 80*time.Millisecond) + st := newPretendStore(t, []byte("x")) + st.stop() + dir := t.TempDir() + var resources []string + for _, m := range []string{"postgres", "keycloak", "minio", "gitea", "mesh-vault", "node-tools"} { + resources = append(resources, `{"id":"`+m+`.bundle-code","type":"archive","source":"`+st.url(m+"/code")+ + `","digest":"`+aDigest+`","path":"`+dir+`/`+m+`"}`) + } + d := declare(t, strings.Join(resources, ",")) + began := time.Now() + _, _, err := ApplyMindingWindows(context.Background(), archHost(t), d, store.State{}, + store.OriginDeclared, noServices, nil, nil, nil, nil) + if err == nil { + t.Fatal("a store that never came back was said fetched") + } + if took := time.Since(began); took > time.Second { + t.Fatalf("six bundles from a store that is down took %s: a wait per bundle", took) + } + if !strings.Contains(err.Error(), "did not answer") { + t.Fatalf("the failure does not say it waited: %v", err) + } + + // An answer is not an absence: 404 fails at once. + patiently(t, 10*time.Second) + missing := newPretendStore(t, []byte("no such blob")) + missing.code = http.StatusNotFound + d = declare(t, `{"id":"x.bundle-code","type":"archive","source":"`+missing.url("x/code")+ + `","digest":"`+aDigest+`","path":"`+dir+`/x"}`) + began = time.Now() + if _, _, err := ApplyMindingWindows(context.Background(), archHost(t), d, store.State{}, + store.OriginDeclared, noServices, nil, nil, nil, nil); err == nil || !strings.Contains(err.Error(), "404") { + t.Fatalf("a missing blob: %v", err) + } + if time.Since(began) > time.Second || missing.fetches() != 1 { + t.Fatalf("a store that answered no was asked %d times", missing.fetches()) + } +} diff --git a/internal/store/lock.go b/internal/store/lock.go index 22245c7..6285faa 100644 --- a/internal/store/lock.go +++ b/internal/store/lock.go @@ -56,3 +56,27 @@ func Lock(statePath string, wait func()) (func(), error) { f.Close() }, nil } + +// TryLockIn takes the apply lock of the state kept in dir when it is free, and never waits: false +// when another apply holds it. For one who must not start while an apply is in flight but must not +// queue behind one either — a scheduled step about to hold a container still (novox/hq issue 291). +func TryLockIn(dir string) (func(), bool, error) { + if err := os.MkdirAll(dir, 0o700); err != nil { + return nil, false, err + } + f, err := os.OpenFile(filepath.Join(dir, LockName), os.O_CREATE|os.O_RDWR, 0o600) + if err != nil { + return nil, false, fmt.Errorf("cannot take this node's apply lock: %w", err) + } + if err := syscall.Flock(int(f.Fd()), syscall.LOCK_EX|syscall.LOCK_NB); err != nil { + f.Close() + if errors.Is(err, syscall.EWOULDBLOCK) { + return nil, false, nil + } + return nil, false, fmt.Errorf("cannot take this node's apply lock: %w", err) + } + return func() { + _ = syscall.Flock(int(f.Fd()), syscall.LOCK_UN) + f.Close() + }, true, nil +}