At 03:30 the store's collector held the registry still while the node-engine's reconcile on that machine was fetching bundle blobs from it: every archive failed 'connection refused' and the machine was held until the next pass. Other machines can meet the same window. A scheduled step now opens its window only once no apply is in flight here (the apply lock is taken just to write the record, so a push still never queues behind the window). An apply whose fetch the store does not answer waits for a window open on its own machine to close, and elsewhere retries with backoff within one bounded budget per apply, well past the window's length; an answer such as 404 still fails at once.
308 lines
11 KiB
Go
308 lines
11 KiB
Go
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())
|
|
}
|
|
}
|