The live tests reached one shared bus and assert, read and remove the mesh's own objects by their fixed names, so packages run in parallel deleted what each other read and the suite passed only one package at a time; a red suite read as noise. internal/testbus starts a server per test, linked in at the nats-server release go.mod pins, and a test holds that pin to the catalogue's bus image and to the facts snapshot's bus when there is one, so the tests never run a bus the mesh does not. The waiter test read a timing (the most connections held at one look) and now reads the state it means (the fewest held across the wait). make check runs the packages in parallel under the race detector, with a timeout.
92 lines
2.5 KiB
Go
92 lines
2.5 KiB
Go
package inventory
|
|
|
|
import (
|
|
"errors"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
// Composing and sending one node's declaration is serialised: a second holder waits for the first.
|
|
func TestHoldingANodeMakesTheNextHolderWait(t *testing.T) {
|
|
inv := fresh(t)
|
|
ctx := t.Context()
|
|
release, err := inv.HoldNodes(ctx, []string{"anchor", "laptop"})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
got := make(chan func(), 1)
|
|
go func() {
|
|
second, err := inv.HoldNodes(ctx, []string{"laptop"})
|
|
if err != nil {
|
|
t.Error(err)
|
|
got <- func() {}
|
|
return
|
|
}
|
|
got <- second
|
|
}()
|
|
select {
|
|
case <-got:
|
|
t.Fatal("a node held by one caller was held by another at the same time")
|
|
case <-time.After(300 * time.Millisecond):
|
|
}
|
|
// Another node is not held up.
|
|
other, err := inv.HoldNodes(ctx, []string{"joiner"})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
other()
|
|
release()
|
|
select {
|
|
case second := <-got:
|
|
second()
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("releasing the node did not let the next holder in")
|
|
}
|
|
}
|
|
|
|
// A waiter pins no pool connection while it waits, and gives up after a bounded wait saying which
|
|
// node is busy.
|
|
func TestAWaiterPinsNoConnectionAndGivesUp(t *testing.T) {
|
|
inv := fresh(t)
|
|
ctx := t.Context()
|
|
release, err := inv.HoldNodes(ctx, []string{"anchor"})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer release()
|
|
savedWait, savedPoll := HoldWaitFor, HoldPoll
|
|
HoldWaitFor, HoldPoll = 1500*time.Millisecond, 50*time.Millisecond
|
|
defer func() { HoldWaitFor, HoldPoll = savedWait, savedPoll }()
|
|
|
|
pool := inv.store.Pool()
|
|
base := pool.Stat().AcquiredConns()
|
|
const waiters = 3
|
|
done := make(chan error, waiters)
|
|
for range waiters {
|
|
go func() {
|
|
_, err := inv.HoldNodes(ctx, []string{"anchor"})
|
|
done <- err
|
|
}()
|
|
}
|
|
// While they wait, the pool lends nothing to them for longer than a look: **pinned is held throughout**,
|
|
// so at some moment of the wait fewer than all of them hold one. Asked of the least seen over many
|
|
// looks, not of the most: a look that lands while every waiter is mid-query is a busy machine, not a
|
|
// pinned connection, and on a loaded one (the suite's packages at once, the race detector) it did land.
|
|
least := waiters
|
|
for range 100 {
|
|
time.Sleep(10 * time.Millisecond)
|
|
if n := int(pool.Stat().AcquiredConns() - base); n < least {
|
|
least = n
|
|
}
|
|
}
|
|
if least >= waiters {
|
|
t.Fatalf("%d waiters held %d connections at every look across their wait", waiters, least)
|
|
}
|
|
for range waiters {
|
|
if err := <-done; !errors.Is(err, ErrNodeBusy) || !strings.Contains(err.Error(), "anchor") {
|
|
t.Fatalf("a waiter did not give up naming the busy node: %v", err)
|
|
}
|
|
}
|
|
}
|