diff --git a/lansport/lansport_test.go b/lansport/lansport_test.go index f62b583..0077281 100644 --- a/lansport/lansport_test.go +++ b/lansport/lansport_test.go @@ -27,12 +27,12 @@ import ( // it is virtual, so it costs nothing to wait out. const probeInterval = time.Second -// rounds lets n probe rounds happen and everything they trigger settle. -func rounds(n int) { - for range n { - time.Sleep(probeInterval) - synctest.Wait() - } +// sleepAndWait advances the bubble's clock by d and then waits for every +// goroutine the passage of time woke to finish what it was doing, so the +// caller can assert on the resulting state. Callers say what d is for. +func sleepAndWait(d time.Duration) { + time.Sleep(d) + synctest.Wait() } var ( @@ -80,13 +80,15 @@ func newStaticPair(t testing.TB, lan *lansporttest.LAN) (a, b *Server) { advA, advB := net.JoinHostPort(ipA.String(), "7890"), net.JoinHostPort(ipB.String(), "7890") a = newServer(t, lan, "a", ipA, advA, StaticSource("a", advB)) b = newServer(t, lan, "b", ipB, advB, StaticSource("b", advA)) - // Round 0 runs at Listen: a can't reach b's advert yet (b didn't - // exist), b fetches a's advert but a doesn't know b's pin. Round 1: a - // fetches b's advert and its LAN check passes, since b knows a. b's - // LAN check races a's fetch, so allow one more round for it. - rounds(2) + // Each server probes once at Listen and then every probeInterval. At + // Listen, a can't reach b's advert (b doesn't exist yet) and b fetches + // a's advert but a doesn't yet know b's pin, so neither is reachable. + // One interval later a fetches b's advert and its LAN check passes, + // since b knows a; b's LAN check of a races a's fetch, so it may need + // the interval after that. Two intervals is therefore enough for both. + sleepAndWait(2 * probeInterval) if len(a.Peers()) != 1 || len(b.Peers()) != 1 { - t.Fatalf("after two rounds: a sees %v, b sees %v", a.Peers(), b.Peers()) + t.Fatalf("a sees %v, b sees %v; want one peer each", a.Peers(), b.Peers()) } return a, b } @@ -144,9 +146,9 @@ func TestStaticPair(t *testing.T) { } c.CloseIdleConnections() - // A peer going away drops out at the next round. + // A peer going away drops out at a's next probe, one interval on. b.Close() - rounds(1) + sleepAndWait(probeInterval) if len(a.Peers()) != 0 { t.Errorf("a still sees %v after b closed", a.Peers()) } @@ -167,7 +169,7 @@ func TestRestartRepins(t *testing.T) { oldTransport := a.Peers()["b"].Transport advB, tlsB := b.AdvertAddr().String(), b.TLSAddr().String() b.Close() - rounds(1) + sleepAndWait(probeInterval) // for a's next probe to drop b if len(a.Peers()) != 0 { t.Fatalf("a still sees %v after b closed", a.Peers()) } @@ -187,7 +189,9 @@ func TestRestartRepins(t *testing.T) { t.Fatal(err) } defer b2.Close() - rounds(2) + // Two intervals for a and the new b to reach each other, for the + // same reasons as in newStaticPair. + sleepAndWait(2 * probeInterval) pb, ok := a.Peers()["b"] if !ok { t.Fatalf("a doesn't see the restarted b: %v", a.Peers()) @@ -281,7 +285,10 @@ func TestTailnetSource(t *testing.T) { } a := mk("a", tsA, ipA) b := mk("b", tsB, ipB) - rounds(2) + // Two intervals for a and b to reach each other, for the same + // reasons as in newStaticPair; they learn of each other from the + // fake tailscaleds at start, before the first probe. + sleepAndWait(2 * probeInterval) pb, ok := a.Peers()["stable-2"] if !ok || pb.Name != "b" || pb.Addr != b.TLSAddr().String() { t.Fatalf("a's peers = %+v", a.Peers()) diff --git a/rendezvous/router.go b/rendezvous/router.go index 289ceda..14cd1c1 100644 --- a/rendezvous/router.go +++ b/rendezvous/router.go @@ -6,33 +6,71 @@ package rendezvous import ( "context" "sync/atomic" + "time" "github.com/tailscale/tb/lansport" ) +// RouterConfig configures a [Router]'s hysteresis: how long membership +// changes must persist before ownership moves. Without it, a peer that flaps +// or a pool whose members are discovered one by one at startup would move +// keys around on every change. +type RouterConfig struct { + // JoinDelay is how long after the most recent newly reachable peer + // appeared the Router waits before giving the new peers ownership of + // keys. Peers that appear within the window are adopted together, so a + // pool coming up one member at a time changes ownership once rather + // than once per member. Until adopted, a new peer owns nothing. Zero + // adopts new peers at once. + JoinDelay time.Duration + + // LeaveDelay is how long a member must stay unreachable before its keys + // are reassigned to other members. Until then it still owns its keys, + // and [Router.Pick] reports them as the local process's to handle, so + // a brief outage or restart neither moves keys away nor back. Zero + // reassigns at once. + LeaveDelay time.Duration +} + // Router picks which of a lansport.Server's peers, or the local process, -// should handle a key. It follows the server's reachable set, so a peer that -// goes away stops being picked within one probe round, and includes the -// local process with its own advertised weight. +// should handle a key. It follows the server's reachable set, with the +// hysteresis in [RouterConfig], and includes the local process with its own +// advertised weight. type Router struct { s *lansport.Server + cfg RouterConfig cur atomic.Pointer[routing] cancel context.CancelFunc done chan struct{} + + // The following fields are only touched by the run goroutine (and by + // NewRouter before it starts), so they need no lock. + adopted map[lansport.Key]float64 // members in the table, by weight; not including the local process + lostAt map[lansport.Key]time.Time // when each adopted member became unreachable, while it is + firstSeen map[lansport.Key]time.Time // when each reachable but not yet adopted peer appeared + timer *time.Timer // for the next adoption or reassignment; nil if none pending } -// routing is one consistent view of the reachable peers and the table built -// from them. Router replaces it wholesale on each change, so a Pick never -// sees a table and a peer map from different rounds. +// routing is one consistent view of the table and the reachable peers. +// Router replaces it wholesale on each change, so a Pick never sees a table +// and a peer map from different rounds. type routing struct { table *Table peers map[lansport.Key]lansport.Peer } // NewRouter returns a Router following s. Close it when done. -func NewRouter(s *lansport.Server) *Router { +func NewRouter(s *lansport.Server, cfg RouterConfig) *Router { ctx, cancel := context.WithCancel(context.Background()) - r := &Router{s: s, cancel: cancel, done: make(chan struct{})} + r := &Router{ + s: s, + cfg: cfg, + cancel: cancel, + done: make(chan struct{}), + adopted: make(map[lansport.Key]float64), + lostAt: make(map[lansport.Key]time.Time), + firstSeen: make(map[lansport.Key]time.Time), + } changes := s.Subscribe() r.update() go r.run(ctx, changes) @@ -48,34 +86,111 @@ func (r *Router) Close() { func (r *Router) run(ctx context.Context, changes <-chan struct{}) { defer close(r.done) for { + var fire <-chan time.Time + if r.timer != nil { + fire = r.timer.C + } select { case <-ctx.Done(): + if r.timer != nil { + r.timer.Stop() + } return case <-changes: - r.update() + case <-fire: + r.timer = nil } + r.update() } } -// update rebuilds the table from the server's reachable peers and itself. +// update brings the adopted member set in line with the server's reachable +// peers, applying the configured delays, rebuilds the table, and arms the +// timer for the next delayed change, if any. func (r *Router) update() { + now := time.Now() peers := r.s.Peers() - members := make(map[string]float64, len(peers)+1) - for key, p := range peers { - members[string(key)] = p.Weight + var next time.Time // when the next delayed change is due; zero if none + + // Reachable peers not yet adopted are joining. Note when each + // appeared, and forget those that left before being adopted. + for key := range peers { + if _, ok := r.adopted[key]; ok { + continue + } + if _, ok := r.firstSeen[key]; !ok { + r.firstSeen[key] = now + } + } + for key := range r.firstSeen { + if _, ok := peers[key]; !ok { + delete(r.firstSeen, key) + } + } + // Adopt all joiners together once JoinDelay has passed since the + // latest one appeared. + if len(r.firstSeen) > 0 { + var latest time.Time + for _, t := range r.firstSeen { + if t.After(latest) { + latest = t + } + } + if due := latest.Add(r.cfg.JoinDelay); !now.Before(due) { + for key := range r.firstSeen { + r.adopted[key] = peers[key].Weight + delete(r.firstSeen, key) + } + } else { + next = due + } + } + + // Adopted members that are unreachable are leaving; drop each one + // LeaveDelay after it was lost, unless it came back first. + for key := range r.adopted { + if p, ok := peers[key]; ok { + r.adopted[key] = p.Weight + delete(r.lostAt, key) + continue + } + lost, ok := r.lostAt[key] + if !ok { + lost = now + r.lostAt[key] = now + } + if due := lost.Add(r.cfg.LeaveDelay); !now.Before(due) { + delete(r.adopted, key) + delete(r.lostAt, key) + } else if next.IsZero() || due.Before(next) { + next = due + } + } + + members := make(map[string]float64, len(r.adopted)+1) + for key, w := range r.adopted { + members[string(key)] = w } - self := r.s.Self() - if self.Key != "" { + if self := r.s.Self(); self.Key != "" { members[string(self.Key)] = self.Weight } table := new(Table) table.Set(members) r.cur.Store(&routing{table: table, peers: peers}) + + if r.timer != nil { + r.timer.Stop() + r.timer = nil + } + if !next.IsZero() { + r.timer = time.NewTimer(next.Sub(now)) + } } // Pick returns the reachable peer that should handle key. It returns false -// when the local process should: because it is the owner, or because there -// are no reachable peers. +// when the local process should: because it is the owner, because the owner +// is unreachable but still within its LeaveDelay, or because there are no +// members besides the local process. func (r *Router) Pick(key string) (lansport.Peer, bool) { cur := r.cur.Load() owner, ok := cur.table.Pick(key) diff --git a/rendezvous/router_test.go b/rendezvous/router_test.go new file mode 100644 index 0000000..e03f48c --- /dev/null +++ b/rendezvous/router_test.go @@ -0,0 +1,186 @@ +// Copyright (c) Tailscale Inc & AUTHORS +// SPDX-License-Identifier: BSD-3-Clause + +package rendezvous + +import ( + "fmt" + "net" + "net/http" + "net/netip" + "slices" + "testing" + "testing/synctest" + "time" + + "github.com/tailscale/tb/lansport" + "github.com/tailscale/tb/lansport/lansporttest" +) + +// probeInterval is the lansport probe interval the tests configure. Under +// synctest it is virtual. +const probeInterval = time.Second + +// sleepAndWait advances the bubble's clock by d and then waits for every +// goroutine the passage of time woke to finish what it was doing, so the +// caller can assert on the resulting state. Callers say what d is for. +func sleepAndWait(d time.Duration) { + time.Sleep(d) + synctest.Wait() +} + +// newServer starts a lansport server named name at ip on lan whose static +// peers are at the given advert addresses, and closes it when the test ends. +func newServer(t testing.TB, lan *lansporttest.LAN, name string, ip netip.Addr, peerAdverts ...string) *lansport.Server { + t.Helper() + s, err := lansport.Listen(lansport.Config{ + Source: lansport.StaticSource(name, peerAdverts...), + Name: name, + AdvertListener: lan.Listen(net.JoinHostPort(ip.String(), "7890")), + TLSListener: lan.Listen(net.JoinHostPort(ip.String(), "0")), + LANIP: ip, + Dial: lan.DialFrom(ip), + Handler: http.NotFoundHandler(), + ProbeInterval: probeInterval, + Logf: func(f string, args ...any) { t.Logf(name+": "+f, args...) }, + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { s.Close() }) + return s +} + +// keysOwnedBy returns up to n keys that r assigns to the member named owner, +// or reports as the local process's if owner is "". +func keysOwnedBy(r *Router, owner string, n int) []string { + var out []string + for i := 0; len(out) < n && i < 10000; i++ { + key := fmt.Sprintf("key-%d", i) + p, ok := r.Pick(key) + if (owner == "" && !ok) || (ok && p.Name == owner) { + out = append(out, key) + } + } + return out +} + +// TestRouterHysteresis checks the join and leave delays: a newly reachable +// peer owns nothing until JoinDelay has passed, a peer that goes away keeps +// its keys (handled locally meanwhile) until LeaveDelay has passed, and one +// that comes back within LeaveDelay resumes as if nothing happened. +func TestRouterHysteresis(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + const joinDelay, leaveDelay = 5 * time.Second, 30 * time.Second + lan := new(lansporttest.LAN) + ipA, ipB := netip.MustParseAddr("127.0.0.1"), netip.MustParseAddr("127.0.0.2") + advA, advB := "127.0.0.1:7890", "127.0.0.2:7890" + a := newServer(t, lan, "a", ipA, advB) + // Let a's probe at start run before b exists, so that a first + // reaches b at its next probe, one interval in, rather than at + // start or one interval in depending on goroutine scheduling. The + // timing assertions below count from there. + synctest.Wait() + b := newServer(t, lan, "b", ipB, advA) + r := NewRouter(a, RouterConfig{JoinDelay: joinDelay, LeaveDelay: leaveDelay}) + defer r.Close() + + // a reaches b at its first probe after b started, one interval in + // (see the lansport tests for why); b takes one more to reach a. b + // is then reachable but not yet adopted. + sleepAndWait(2 * probeInterval) + if len(a.Peers()) != 1 { + t.Fatalf("a's peers = %v, want b", a.Peers()) + } + if got := r.Members(); !slices.Equal(got, []string{"a"}) { + t.Fatalf("members right after b appeared = %v, want [a]", got) + } + if got := keysOwnedBy(r, "b", 1); len(got) != 0 { + t.Errorf("b owns %v before adoption", got) + } + + // Adoption is due JoinDelay after a first saw b, which was one + // interval in. One interval short of that, still not adopted; at + // it, adopted. + sleepAndWait(joinDelay - 2*probeInterval) + if got := r.Members(); !slices.Equal(got, []string{"a"}) { + t.Fatalf("members before JoinDelay = %v, want [a]", got) + } + sleepAndWait(probeInterval) + if got := r.Members(); !slices.Equal(got, []string{"a", "b"}) { + t.Fatalf("members after JoinDelay = %v, want [a b]", got) + } + bKeys := keysOwnedBy(r, "b", 3) + if len(bKeys) != 3 { + t.Fatalf("b owns no keys after adoption") + } + + // b goes away. One probe interval later a has noticed; b keeps + // its keys, which the local process handles meanwhile. + b.Close() + sleepAndWait(probeInterval) + if len(a.Peers()) != 0 { + t.Fatalf("a still sees %v after b closed", a.Peers()) + } + if got := r.Members(); !slices.Equal(got, []string{"a", "b"}) { + t.Fatalf("members right after b left = %v, want [a b] during LeaveDelay", got) + } + for _, key := range bKeys { + if p, ok := r.Pick(key); ok { + t.Errorf("Pick(%s) during b's LeaveDelay = %v, want local", key, p.Name) + } + } + + // b comes back (same name, so the same routing key) well within + // LeaveDelay: nothing moved, and its keys route to it again as + // soon as it is reachable, with no JoinDelay since it was never + // dropped. + b = newServer(t, lan, "b", ipB, advA) + sleepAndWait(2 * probeInterval) + if got := r.Members(); !slices.Equal(got, []string{"a", "b"}) { + t.Fatalf("members after b returned = %v, want [a b]", got) + } + for _, key := range bKeys { + if p, ok := r.Pick(key); !ok || p.Name != "b" { + t.Errorf("Pick(%s) after b returned = %v, %v; want b", key, p.Name, ok) + } + } + + // b goes away for good: LeaveDelay after a notices, its keys are + // reassigned, which with two members means they are a's. + b.Close() + sleepAndWait(probeInterval) + sleepAndWait(leaveDelay - probeInterval) + if got := r.Members(); !slices.Equal(got, []string{"a", "b"}) { + t.Fatalf("members just before LeaveDelay = %v, want [a b]", got) + } + sleepAndWait(probeInterval) + if got := r.Members(); !slices.Equal(got, []string{"a"}) { + t.Fatalf("members after LeaveDelay = %v, want [a]", got) + } + }) +} + +// TestRouterNoDelays checks that zero delays mean immediate adoption and +// reassignment, which is what callers that want the old behavior, and +// tests, get. +func TestRouterNoDelays(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + lan := new(lansporttest.LAN) + ipA, ipB := netip.MustParseAddr("127.0.0.1"), netip.MustParseAddr("127.0.0.2") + a := newServer(t, lan, "a", ipA, "127.0.0.2:7890") + b := newServer(t, lan, "b", ipB, "127.0.0.1:7890") + r := NewRouter(a, RouterConfig{}) + defer r.Close() + + sleepAndWait(2 * probeInterval) // for a and b to reach each other + if got := r.Members(); !slices.Equal(got, []string{"a", "b"}) { + t.Fatalf("members = %v, want [a b]", got) + } + b.Close() + sleepAndWait(probeInterval) // for a's next probe to drop b + if got := r.Members(); !slices.Equal(got, []string{"a"}) { + t.Fatalf("members after b closed = %v, want [a]", got) + } + }) +}