package lease import ( "context" "errors" "fmt" "os" "sync" "testing" "time" "github.com/nats-io/nats.go" "github.com/nats-io/nats.go/jetstream" ) // Two controllers at once, on a real bus (novox/hq to-be 45 §6, replay R1's half that lives here): // one takes the lease and the other waits; a handover gives the second a higher epoch; a holder that // stops renewing is taken over once its key expires, and has stopped acting before then; a renewal // refused is a loss at once; and a bucket raised again from nothing never issues an epoch twice. // // MESH_TEST_NATS=nats://127.0.0.1:14222 go test ./internal/lease/ // The bounds in these tests are the design's divided by five, so the cases run in seconds. const ( testTTL = 3 * time.Second testRenew = time.Second testPoll = 100 * time.Millisecond ) // aBucket is a fresh lease bucket on the test bus, with the test's age. func aBucket(t *testing.T) (jetstream.JetStream, string) { t.Helper() url := os.Getenv("MESH_TEST_NATS") if url == "" { t.Skip("MESH_TEST_NATS unset") } conn, err := nats.Connect(url) if err != nil { t.Fatal(err) } t.Cleanup(conn.Close) js, err := jetstream.New(conn) if err != nil { t.Fatal(err) } bucket := fmt.Sprintf("lease-test-%d", time.Now().UnixNano()) if _, err := js.CreateKeyValue(t.Context(), jetstream.KeyValueConfig{Bucket: bucket, History: 1, TTL: testTTL, Storage: jetstream.FileStorage}); err != nil { t.Fatal(err) } t.Cleanup(func() { _ = js.DeleteKeyValue(context.Background(), bucket) }) return js, bucket } // aController is one controller instance's lease over the bucket, on a connection of its own. func aController(t *testing.T, bucket, name string, floor func(context.Context) (uint64, error)) *Lease { t.Helper() conn, err := nats.Connect(os.Getenv("MESH_TEST_NATS")) if err != nil { t.Fatal(err) } t.Cleanup(conn.Close) js, err := jetstream.New(conn) if err != nil { t.Fatal(err) } var said sync.Mutex l, err := Open(t.Context(), js, bucket, Options{Holder: Holder{Instance: name, Host: name}, RenewEvery: testRenew, Margin: testRenew / 2, Poll: testPoll, Floor: floor, Say: func(format string, args ...any) { said.Lock() defer said.Unlock() t.Logf(name+": "+format, args...) }}) if err != nil { t.Fatal(err) } return l } // takeWithin takes the lease in a goroutine and answers when it did, or fails past the bound. func takeWithin(t *testing.T, l *Lease) <-chan uint64 { t.Helper() took := make(chan uint64, 1) go func() { epoch, err := l.Take(t.Context()) if err != nil { t.Errorf("taking: %v", err) close(took) return } took <- epoch }() return took } func TestTheSecondControllerWaitsAndTakesAHigherEpochOnHandover(t *testing.T) { _, bucket := aBucket(t) a, b := aController(t, bucket, "a", nil), aController(t, bucket, "b", nil) epochA, err := a.Take(t.Context()) if err != nil { t.Fatal(err) } keeping, stop := context.WithCancel(t.Context()) defer stop() go a.Keep(keeping) // b waits while a renews — across more than one age of the key, so it is the renewals holding it. took := takeWithin(t, b) select { case epoch := <-took: t.Fatalf("b took the lease (epoch %d) while a held it and renewed", epoch) case <-time.After(testTTL + testRenew): } if _, err := b.Epoch(); !errors.Is(err, ErrNotHeld) { t.Fatalf("b, waiting, may act: %v", err) } if e, err := a.Epoch(); err != nil || e != epochA { t.Fatalf("a, holding, answers %d, %v", e, err) } holder, found, err := b.Current(t.Context()) if err != nil || !found || holder.Instance != "a" || holder.Epoch != epochA { t.Fatalf("b reads the holder as %+v (%v, %v), want a at epoch %d", holder, found, err, epochA) } // a gives it back: b takes it at once, not after the key's age, at a higher epoch. stop() handedOver := time.Now() a.Release(t.Context()) select { case epochB := <-took: if epochB <= epochA { t.Fatalf("b took epoch %d after a's %d: an epoch must only grow", epochB, epochA) } if waited := time.Since(handedOver); waited > testTTL { t.Fatalf("b took the lease %s after a gave it back: it waited out the key's age", waited) } case <-time.After(2 * testTTL): t.Fatal("b never took the lease a gave back") } if _, err := a.Epoch(); !errors.Is(err, ErrNotHeld) { t.Fatalf("a, having given it back, may still act: %v", err) } if _, err := a.TryTake(t.Context()); err == nil { t.Fatal("a took the lease again after giving it back: its successor is a new candidate") } } func TestAHolderThatStopsRenewingStopsActingBeforeItIsTakenOver(t *testing.T) { _, bucket := aBucket(t) a, b := aController(t, bucket, "a", nil), aController(t, bucket, "b", nil) epochA, err := a.Take(t.Context()) if err != nil { t.Fatal(err) } // a is stuck: it never renews (a process stopped, a goroutine wedged). It was not told anything. took := takeWithin(t, b) // Its gate closes by its own clock, before the key can expire. closed := time.Now() for a.Held() { if time.Since(closed) > testTTL { t.Fatal("a still may act past its key's age without renewing") } time.Sleep(20 * time.Millisecond) } select { case epoch := <-took: t.Fatalf("b took the lease (epoch %d) before a stopped acting", epoch) default: } select { case epochB := <-took: if epochB <= epochA { t.Fatalf("b took epoch %d after a's %d", epochB, epochA) } // And a, the moment it tries to renew, knows it lost: the key moved. if err := a.Renew(t.Context()); err == nil { t.Fatal("a renewed a lease b holds") } select { case <-a.Lost(): default: t.Fatal("a's renewal was refused and a was not told it lost the lease") } if _, err := a.Epoch(); !errors.Is(err, ErrNotHeld) { t.Fatalf("a, having lost the lease, may act: %v", err) } case <-time.After(3 * testTTL): t.Fatal("b never took over a lease nobody renewed") } } func TestARenewalRefusedIsALossAtOnce(t *testing.T) { js, bucket := aBucket(t) a := aController(t, bucket, "a", nil) if _, err := a.Take(t.Context()); err != nil { t.Fatal(err) } // Somebody else writes the key — a second controller that read it as expired on a skewed clock, // or a person — so a's next renewal finds a revision it did not write. kv, err := js.KeyValue(t.Context(), bucket) if err != nil { t.Fatal(err) } if _, err := kv.Put(t.Context(), Key, []byte(`{"instance":"intruder"}`)); err != nil { t.Fatal(err) } if err := a.Renew(t.Context()); err == nil { t.Fatal("a renewed over a key somebody else wrote") } select { case <-a.Lost(): case <-time.After(time.Second): t.Fatal("a was not told it lost the lease") } if _, err := a.Epoch(); !errors.Is(err, ErrNotHeld) { t.Fatalf("a acts after its renewal was refused: %v", err) } // And giving back a lease it lost deletes nothing of the one who holds it now. a.Release(t.Context()) if h, found, err := Current(t.Context(), kv); err != nil || !found || h.Instance != "intruder" { t.Fatalf("after a gave back what it lost, the key holds %+v (%v, %v)", h, found, err) } } func TestKeepingRenewsAcrossManyAges(t *testing.T) { _, bucket := aBucket(t) a := aController(t, bucket, "a", nil) epoch, err := a.Take(t.Context()) if err != nil { t.Fatal(err) } keeping, stop := context.WithCancel(t.Context()) defer stop() go a.Keep(keeping) until := time.Now().Add(3 * testTTL) for time.Now().Before(until) { if e, err := a.Epoch(); err != nil || e != epoch { t.Fatalf("a, renewing, answers %d, %v", e, err) } time.Sleep(200 * time.Millisecond) } h, found, err := a.Current(t.Context()) if err != nil || !found || h.Epoch != epoch || h.Instance != "a" { t.Fatalf("the key says %+v (%v, %v)", h, found, err) } } func TestAnEpochIsNeverIssuedTwiceWhenTheBucketStartsOver(t *testing.T) { _, bucket := aBucket(t) // The mesh has issued epoch 500 before: its store says so. The bucket is new — a bus whose data // was replaced — and its revisions start at one. floor := func(context.Context) (uint64, error) { return 500, nil } a := aController(t, bucket, "a", floor) moved := false a.o.Moved = func(was, floor uint64) { moved = was < floor && floor == 500 } epoch, err := a.Take(t.Context()) if err != nil { t.Fatal(err) } if epoch <= 500 { t.Fatalf("a took epoch %d with 500 already issued: every machine that heard 500 would refuse it", epoch) } if !moved { t.Fatal("the bucket's revisions were moved past the floor and nobody was told") } // A floor that cannot be read takes nothing: an epoch that may not be higher is not issued. _, other := aBucket(t) b := aController(t, other, "b", func(context.Context) (uint64, error) { return 0, errors.New("store away") }) if _, err := b.TryTake(t.Context()); err == nil { t.Fatal("b took a lease without knowing the highest epoch issued") } } func TestALeaseBucketWithoutAnAgeIsRefused(t *testing.T) { url := os.Getenv("MESH_TEST_NATS") if url == "" { t.Skip("MESH_TEST_NATS unset") } conn, err := nats.Connect(url) if err != nil { t.Fatal(err) } defer conn.Close() js, _ := jetstream.New(conn) bucket := fmt.Sprintf("lease-ageless-%d", time.Now().UnixNano()) if _, err := js.CreateKeyValue(t.Context(), jetstream.KeyValueConfig{Bucket: bucket, History: 1}); err != nil { t.Fatal(err) } defer func() { _ = js.DeleteKeyValue(context.Background(), bucket) }() if _, err := Open(t.Context(), js, bucket, Options{Holder: Holder{Instance: "a"}}); err == nil { t.Fatal("a lease over keys that never expire was opened: a controller that died holding it would hold it for ever") } }