Files
mesh-controller/internal/link/store_window_test.go
T

182 lines
6.7 KiB
Go

package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"testing"
"github.com/jackc/pgx/v5/pgconn"
amqp "github.com/rabbitmq/amqp091-go"
)
// What a store restarting under an adoption answers with (novox/hq issues 082, 083).
var restarting = &pgconn.PgError{Code: "57P03", Message: "the database system is starting up"}
type recordsWith struct{ err error }
func (r recordsWith) Built(context.Context, BuildResult) error { return r.err }
type upgradesWith struct{ err error }
func (u upgradesWith) Upgraded(context.Context, Upgraded) error { return u.err }
type replaysWith struct{ err error }
func (r replaysWith) Announceable(context.Context) ([]Announcement, error) { return nil, r.err }
type settledAs struct{ acked, nacked, requeued, rejected bool }
func (a *settledAs) Ack(uint64, bool) error { a.acked = true; return nil }
func (a *settledAs) Nack(_ uint64, _ bool, requeue bool) error {
a.nacked, a.requeued = true, requeue
return nil
}
func (a *settledAs) Reject(uint64, bool) error { a.rejected = true; return nil }
func a(t *testing.T, to *settledAs, key string, v any) amqp.Delivery {
t.Helper()
body, err := json.Marshal(v)
if err != nil {
t.Fatal(err)
}
tag++
return amqp.Delivery{Acknowledger: to, RoutingKey: key, Body: body, DeliveryTag: tag}
}
func quietServer() *Server { return &Server{log: log.New(io.Discard, "", 0)} }
func (a *settledAs) held() bool { return !a.acked && !a.nacked && !a.rejected }
// A build result the store could not take right now is handed back; one it refused is rejected,
// as before; one it kept is acknowledged.
func TestABuildResultWaitsOutARestartingStore(t *testing.T) {
built := BuildResult{On: "anchor", Repository: "/r", Commit: "abc"}
for _, c := range []struct {
what string
err error
want func(*settledAs) bool
}{
{"restarting", restarting, func(s *settledAs) bool { return s.held() }},
{"refused", errors.New("no such module"), func(s *settledAs) bool { return s.rejected && !s.nacked }},
{"kept", nil, func(s *settledAs) bool { return s.acked && !s.nacked }},
} {
s := quietServer()
s.recorder = recordsWith{err: c.err}
to := &settledAs{}
s.handleBuilt(context.Background(), a(t, to, KeyBuilt, built))
if !c.want(to) {
t.Errorf("%s: a build result was settled as %+v", c.what, *to)
}
}
}
// An upgrade announcement arriving while the store restarts is asked again; any other failure is
// acknowledged, so it cannot stop every upgrade behind it.
func TestAnUpgradeWaitsOutARestartingStoreAndNothingElse(t *testing.T) {
moved := Upgraded{Module: "gitea", Commit: "abcdef0123"}
for _, c := range []struct {
what string
err error
want func(*settledAs) bool
}{
{"the store away, said by the upgrader", errors.Join(ErrTryAgain, restarting), func(s *settledAs) bool { return s.held() }},
{"a push that timed out", context.DeadlineExceeded, func(s *settledAs) bool { return s.acked && !s.nacked }},
{"cannot act", errors.New("anchor cannot be resolved"), func(s *settledAs) bool { return s.acked && !s.nacked }},
{"acted", nil, func(s *settledAs) bool { return s.acked && !s.nacked }},
} {
s := quietServer()
s.upgrader = upgradesWith{err: c.err}
to := &settledAs{}
s.upgraded(context.Background(), a(t, to, "upgraded", moved))
if !c.want(to) {
t.Errorf("%s: an upgrade was settled as %+v", c.what, *to)
}
}
}
// A catalogue's request to catch up is acknowledged after the work, and asked again while the
// store cannot be read — not lost until the catalogue next restarts.
func TestACatchUpWaitsOutARestartingStore(t *testing.T) {
for _, c := range []struct {
what string
err error
want func(*settledAs) bool
}{
{"restarting", restarting, func(s *settledAs) bool { return s.held() }},
{"unreadable", errors.New("a build row is malformed"), func(s *settledAs) bool { return s.acked && !s.nacked }},
{"nothing to replay", nil, func(s *settledAs) bool { return s.acked && !s.nacked }},
} {
s := quietServer()
s.replayer = replaysWith{err: c.err}
to := &settledAs{}
s.catchingUp(context.Background(), a(t, to, "catch-up", map[string]string{}))
if !c.want(to) {
t.Errorf("%s: a catch-up request was settled as %+v", c.what, *to)
}
}
}
// Shutting down is not an answer about a message: one handled with a cancelled context is left
// unsettled, for the broker to hand to whatever consumes next (issue 083, review).
func TestAMessageHandledDuringShutdownIsLeftForTheBroker(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel()
s := quietServer()
s.recorder = recordsWith{err: context.Canceled}
to := &settledAs{}
s.handleBuilt(ctx, a(t, to, KeyBuilt, BuildResult{On: "anchor", Repository: "/r", Commit: "abc"}))
if !to.held() {
t.Fatalf("a build result handled during shutdown was settled, and so lost: %+v", *to)
}
}
// Two identical build results: the newer sets the older aside rather than leaving it unsettled
// for ever, holding a place in the prefetch.
func TestAnIdenticalBuildResultSetsTheHeldOneAside(t *testing.T) {
s := quietServer()
s.recorder = recordsWith{err: restarting}
built := BuildResult{On: "anchor", Repository: "/r", Commit: "abc"}
first, second := &settledAs{}, &settledAs{}
s.handleBuilt(context.Background(), a(t, first, KeyBuilt, built))
s.handleBuilt(context.Background(), a(t, second, KeyBuilt, built))
if !first.acked || !second.held() || len(s.parked) != 1 {
t.Fatalf("an identical build result did not set the held one aside: first %+v second %+v, %d held",
*first, *second, len(s.parked))
}
}
// What is held stops short of the prefetch, so the loop always has room to answer an enrolment.
func TestWhatIsHeldLeavesRoomInThePrefetch(t *testing.T) {
s := quietServer()
s.recorder = recordsWith{err: restarting}
var last *settledAs
for i := 0; i < Prefetch; i++ {
last = &settledAs{}
s.handleBuilt(context.Background(), a(t, last, KeyBuilt, BuildResult{On: "anchor", Commit: fmt.Sprint(i)}))
}
if len(s.parked) != Prefetch-PrefetchHeadroom {
t.Fatalf("%d messages were held; the ceiling is %d", len(s.parked), Prefetch-PrefetchHeadroom)
}
if last.held() {
t.Fatalf("a message past the ceiling was held: %+v", *last)
}
}
// An upgrade handled during shutdown is left for the broker too — the upgrader's error is the
// cancelled context, which is no answer about the announcement.
func TestAnUpgradeHandledDuringShutdownIsLeftForTheBroker(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel()
s := quietServer()
s.upgrader = upgradesWith{err: context.Canceled}
to := &settledAs{}
s.upgraded(ctx, a(t, to, "upgraded", Upgraded{Module: "gitea", Commit: "abcdef0123"}))
if !to.held() {
t.Fatalf("an upgrade handled during shutdown was settled, and so lost: %+v", *to)
}
}