Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 24 additions & 17 deletions lansport/lansport_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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())
}
Expand All @@ -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())
}
Expand All @@ -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())
Expand Down Expand Up @@ -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())
Expand Down
149 changes: 132 additions & 17 deletions rendezvous/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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)
Expand Down
Loading
Loading