mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-09 23:19:11 +02:00
* Skip session warnings that fire after their window The warning timers run on the monotonic clock, which does not advance while an Android device is suspended. A timer armed for T-10 or T-2 can therefore fire long after the window it was armed for, delivering a "session expires soon" notification once that window is already gone. Gate both callbacks on the wall clock at fire time: the T-10 warning is skipped once the final-warning window has been reached, and the final warning is skipped once the deadline itself has passed. Both set their edge guard before returning so a skipped warning cannot fire again for the same deadline. * Harden the late-warning guards Clamp a non-positive final lead to zero in the T-10 guard so a disabled final warning cannot move the cutoff past the deadline, matching how armTimerLocked already treats it. Strip the monotonic reading from both sides of the comparison so the guard measures wall-clock time regardless of how the caller built the deadline. The production deadline comes from a protobuf timestamp and has no monotonic reading; this keeps the guard correct for callers that derive one from time.Now. * Log the deadline and lateness on skipped warnings Include the deadline and how far past the cutoff the timer fired, so a debug bundle shows how long the device was suspended. * Inject the clock into the late-warning guard and cover it with tests The guard read time.Now internally, so the skip paths were reachable only through a deadline already in the past and the boundary depended on real time. Extract the comparison into isLate and read the time through a nowFn field, so tests can place a resume anywhere around the deadline without sleeping. * Send the final warning when the T-10 timer fires inside its window A suspend between roughly eight and ten minutes long made the T-10 timer fire inside the final-warning window and the final timer fire after the deadline, so both were skipped and a user who resumed with time left got no warning at all. When the T-10 timer fires late but before the deadline, send the final warning in its place and mark it fired so the delayed final timer does not repeat it. * Respect dismissal when promoting a late warning to the final one fireFinal skips the final warning once the user dismissed the deadline, but the promoted path did not, so a dismissed deadline could still get a final warning. Check the dismissal first, and give each skip reason its own log line so an already-fired final warning no longer logs a negative lateness. * Add a deadline-only mode to the session watcher Android will schedule its own expiry warnings from the deadline, so the engine must not arm the T-10 and T-2 timers there. NewDeadlineOnly keeps the deadline validation, the recorder propagation and the logging, and skips only the timers, so the status snapshot the app reads stays correct and an out-of-range deadline is still rejected. * Use the deadline-only watcher on Android and drop the warning callbacks The warning timers run on the monotonic clock, which does not advance while the device sleeps, so a warning armed for T-10 could fire long after its window. The app now schedules the warnings itself with WorkManager, anchored to the wall clock, from the deadline it reads through SessionExpiresAtUnix on every OnStateChanged. Wire the deadline-only watcher into the android build and remove the event-driven path from the gomobile surface: OnSessionExpiring, the event subscription behind it and DismissSessionWarning, which the app never called. * Describe the late-warning guard without naming Android The guard stays for the desktop builds, where a timer can also stall across a sleep. Android no longer arms the timers at all. * [client] Evaluate the session deadline against the wall clock The session deadline is an absolute instant published by management, but the warnings for it were armed as relative timers. A relative timer runs on the monotonic clock, which does not advance while a device is suspended: it fires once that much awake time has passed, which can be long after the deadline, and nothing re-evaluates when the device wakes up. A device that suspends before the warning window and resumes with minutes left is never told it can still extend. The same applies to a device that boots before NTP has corrected its clock: the one-shot timer has already fired by the time the correction lands. Compare the tracked deadline against the wall clock on a ticker instead. Every tick after a resume, or after a clock correction, sees the real remaining time, so the warning is published whenever there is still time to act on it and never once the deadline has passed. This supersedes the late-callback guards: with no relative timer there is no late callback to detect. The warning windows keep their semantics. Each one publishes at most once per deadline value, a dismissal still suppresses the final warning, and inside the final window the interactive warning is skipped in favour of the final one, since a device resuming there never saw it. Update evaluates immediately so a deadline that already sits inside a window does not wait for a tick, and Close stops the loop and waits for it. Two trade-offs: a warning can land up to one tick (10s) after its lead, which is noise against leads measured in minutes, and the loop runs from the first tracked deadline until Close rather than being armed and disarmed per deadline. * [client] Poll whatever the clock reads when the deadline arrives Update started the loop only for a deadline in the future, and read the wall clock directly to decide. A device whose clock runs ahead before NTP corrects it accepts the deadline as recent-past, starts nothing, and then has no loop left to notice the correction: the next sync carries the same deadline value, which Update treats as a no-op. The warning is lost for exactly the case this mechanism exists to cover. Start the loop for every accepted deadline. evaluate already ignores a deadline that has genuinely passed, so the only cost is a ticker on an expired session until the next deadline or Close. Take the current time from nowFn there as well, so the sanity checks and the evaluation agree on what time it is and a test can drive both. * [client] Announce the deadline before any warning about it Update released the lock before telling the recorder about the new deadline, so a tick landing in that window could publish a warning for a deadline consumers had not been told about yet, inverting the order the function documents. Gate publishing on the deadline the recorder has been told about: Update records it after the recorder call, and an evaluation that finds the gate shut leaves the warning for the next tick. * [client] Describe the deadline-only mode in terms of the poll The mode no longer skips arming timers, it skips the evaluation entirely, and the doc comment said otherwise. * [client] Publish warnings only from the evaluation loop Update evaluated the new deadline on the caller's goroutine, so a warning could still be on its way to the recorder after Close had returned: Close waits for the evaluation loop, and that publish was not coming from it. Hand the work to the loop instead. Update announces the deadline and nudges a buffered wake channel, so a deadline that already sits inside a warning window is still warned about at once, without this goroutine ever touching the recorder. The loop is now the only caller of evaluate, which makes waiting for it in Close enough. * [client] Make a concurrent Close wait for the loop as well The first Close stopped the evaluation loop and waited for it, but a second concurrent Close saw the closed flag and returned straight away, telling its caller the teardown was done while a warning was still on its way to the recorder. Keep the completion channel on the receiver after the teardown starts, so the call that loses the race waits on the same loop. * [client] Correct the poll start condition in the doc comment Update starts the loop for every accepted deadline, including one that already reads as expired, not only for a future one. * [client] Unblock the parked publish on every test exit path The concurrent-Close test parks the evaluation loop inside a publish and releases it at the end. A t.Fatal before that line runs Goexit and skips the release, stranding the loop and both Close calls for the rest of the package run. Release through a sync.Once registered with t.Cleanup, and close the watcher there too, so an early failure reports itself instead of hanging. * [client] Release the parked publish before closing in the test cleanup t.Cleanup runs in reverse registration order, so the watcher was closed before the release ran: Close waits for the evaluation loop, the loop was still parked in the publish, and a t.Fatal hung the cleanup instead of reporting the failure. Register one cleanup that releases first and closes after. --------- Co-authored-by: Zoltán Papp <zoltan.pmail@gmail.com>
920 lines
27 KiB
Go
920 lines
27 KiB
Go
package sessionwatch
|
|
|
|
import (
|
|
"errors"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
cProto "github.com/netbirdio/netbird/client/proto"
|
|
)
|
|
|
|
// fakeRecorder satisfies StatusRecorder and records every call so tests
|
|
// can observe what the watcher emits. SetSessionExpiresAt and PublishEvent
|
|
// land in the same ordered events slice (with the Kind distinguishing
|
|
// them) so tests that care about ordering still work. lastDeadline holds
|
|
// the most recent value passed to SetSessionExpiresAt so tests can assert
|
|
// the recorder ended up cleared/set as expected.
|
|
type fakeRecorder struct {
|
|
mu sync.Mutex
|
|
events []event
|
|
lastDeadline time.Time
|
|
// setDelay stalls SetSessionExpiresAt, widening the window in which a
|
|
// concurrent evaluation could publish a warning out of order.
|
|
setDelay time.Duration
|
|
}
|
|
|
|
type eventKind int
|
|
|
|
const (
|
|
stateChange eventKind = iota
|
|
publish
|
|
)
|
|
|
|
type event struct {
|
|
kind eventKind
|
|
// Set only for publish events.
|
|
severity cProto.SystemEvent_Severity
|
|
category cProto.SystemEvent_Category
|
|
message string
|
|
meta map[string]string
|
|
}
|
|
|
|
// SetSessionExpiresAt mirrors peer.Status: a same-value write is a no-op,
|
|
// a real change records the new value and fans out a state-change (the
|
|
// production recorder calls notifyStateChange internally). The baseline
|
|
// is the zero time, so an initial clear before any deadline is set emits
|
|
// nothing — matching the real recorder.
|
|
func (r *fakeRecorder) SetSessionExpiresAt(deadline time.Time) {
|
|
r.mu.Lock()
|
|
delay := r.setDelay
|
|
r.mu.Unlock()
|
|
// Stall without the lock held, so a concurrent publish can still record.
|
|
time.Sleep(delay)
|
|
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
if r.lastDeadline.Equal(deadline) {
|
|
return
|
|
}
|
|
r.lastDeadline = deadline
|
|
r.events = append(r.events, event{kind: stateChange})
|
|
}
|
|
|
|
func (r *fakeRecorder) deadline() time.Time {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
return r.lastDeadline
|
|
}
|
|
|
|
func (r *fakeRecorder) PublishEvent(
|
|
severity cProto.SystemEvent_Severity,
|
|
category cProto.SystemEvent_Category,
|
|
message string,
|
|
_ string,
|
|
metadata map[string]string,
|
|
) {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
r.events = append(r.events, event{
|
|
kind: publish,
|
|
severity: severity,
|
|
category: category,
|
|
message: message,
|
|
meta: metadata,
|
|
})
|
|
}
|
|
|
|
func (r *fakeRecorder) snapshot() []event {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
out := make([]event, len(r.events))
|
|
copy(out, r.events)
|
|
return out
|
|
}
|
|
|
|
func (e event) isFinalWarning() bool {
|
|
return e.kind == publish && e.meta[MetaSessionFinal] == "true"
|
|
}
|
|
|
|
func (e event) isWarning() bool {
|
|
return e.kind == publish && e.meta[MetaSessionWarning] == "true" && e.meta[MetaSessionFinal] != "true"
|
|
}
|
|
|
|
func countWhere(events []event, pred func(event) bool) int {
|
|
n := 0
|
|
for _, e := range events {
|
|
if pred(e) {
|
|
n++
|
|
}
|
|
}
|
|
return n
|
|
}
|
|
|
|
func waitForEvents(t *testing.T, r *fakeRecorder, want int) []event {
|
|
t.Helper()
|
|
deadline := time.Now().Add(500 * time.Millisecond)
|
|
for time.Now().Before(deadline) {
|
|
if got := r.snapshot(); len(got) >= want {
|
|
return got
|
|
}
|
|
time.Sleep(5 * time.Millisecond)
|
|
}
|
|
got := r.snapshot()
|
|
t.Fatalf("timed out waiting for %d events, got %d: %+v", want, len(got), got)
|
|
return nil
|
|
}
|
|
|
|
// testInterval keeps the ticker-driven tests fast; the leads they use are
|
|
// in the same millisecond scale.
|
|
const testInterval = 2 * time.Millisecond
|
|
|
|
// newWatcher builds a watcher with the final warning disabled (finalLead=0),
|
|
// matching the lead-only behaviour the pre-final-warning tests assume.
|
|
func newWatcher(lead time.Duration, r *fakeRecorder) *Watcher {
|
|
return newWatcherWithLeads(lead, 0, r)
|
|
}
|
|
|
|
// newWatcherWithLeads builds a watcher that evaluates on testInterval. The
|
|
// interval is set before Update, so the evaluation loop does not exist yet
|
|
// and the write cannot race it.
|
|
func newWatcherWithLeads(lead, final time.Duration, r StatusRecorder) *Watcher {
|
|
w := NewWithLeads(lead, final, r)
|
|
w.interval = testInterval
|
|
return w
|
|
}
|
|
|
|
// fakeClock is the watcher's wall clock under test. Reads come from the
|
|
// evaluation goroutine while the test writes, so both go through the mutex.
|
|
type fakeClock struct {
|
|
mu sync.Mutex
|
|
t time.Time
|
|
}
|
|
|
|
func newFakeClock(t time.Time) *fakeClock {
|
|
return &fakeClock{t: t}
|
|
}
|
|
|
|
func (c *fakeClock) now() time.Time {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
return c.t
|
|
}
|
|
|
|
// set jumps the clock, standing in for a resume from suspension or for an
|
|
// NTP correction.
|
|
func (c *fakeClock) set(t time.Time) {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
c.t = t
|
|
}
|
|
|
|
// settle waits out a handful of evaluation ticks so a test can assert that
|
|
// nothing was published.
|
|
func settle() {
|
|
time.Sleep(20 * testInterval)
|
|
}
|
|
|
|
func TestUpdateZeroBeforeAnythingIsNoop(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(50*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
_ = w.Update(time.Time{})
|
|
|
|
if got := r.snapshot(); len(got) != 0 {
|
|
t.Fatalf("expected no events on initial zero, got %+v", got)
|
|
}
|
|
}
|
|
|
|
func TestUpdateNonZeroFiresStateChange(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(50*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
d := time.Now().Add(time.Hour)
|
|
_ = w.Update(d)
|
|
|
|
events := waitForEvents(t, r, 1)
|
|
if events[0].kind != stateChange {
|
|
t.Fatalf("expected stateChange, got %+v", events[0])
|
|
}
|
|
if !w.Deadline().Equal(d) {
|
|
t.Fatalf("deadline mismatch: %v vs %v", w.Deadline(), d)
|
|
}
|
|
}
|
|
|
|
func TestSameDeadlineIsNoop(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(50*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
d := time.Now().Add(time.Hour)
|
|
_ = w.Update(d)
|
|
_ = w.Update(d)
|
|
_ = w.Update(d)
|
|
|
|
events := waitForEvents(t, r, 1)
|
|
if len(events) != 1 {
|
|
t.Fatalf("expected exactly 1 event for repeated same deadline, got %d: %+v", len(events), events)
|
|
}
|
|
}
|
|
|
|
func TestWarningFiresOnceWithinLeadWindow(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
lead := 50 * time.Millisecond
|
|
w := newWatcher(lead, r)
|
|
defer w.Close()
|
|
|
|
// Deadline 80ms out — warning should fire after ~30ms.
|
|
d := time.Now().Add(80 * time.Millisecond)
|
|
_ = w.Update(d)
|
|
|
|
events := waitForEvents(t, r, 2)
|
|
if events[0].kind != stateChange {
|
|
t.Fatalf("event[0] should be stateChange, got %+v", events[0])
|
|
}
|
|
if !events[1].isWarning() {
|
|
t.Fatalf("event[1] should be a warning publish, got %+v", events[1])
|
|
}
|
|
}
|
|
|
|
func TestWarningFiresImmediatelyWhenAlreadyInsideWindow(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(time.Hour, r) // lead > delta => fire immediately
|
|
defer w.Close()
|
|
|
|
d := time.Now().Add(10 * time.Millisecond)
|
|
_ = w.Update(d)
|
|
|
|
events := waitForEvents(t, r, 2)
|
|
if !events[1].isWarning() {
|
|
t.Fatalf("expected immediate warning publish, got %+v", events[1])
|
|
}
|
|
}
|
|
|
|
func TestNewDeadlineCancelsPriorTimer(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
lead := 50 * time.Millisecond
|
|
w := newWatcher(lead, r)
|
|
defer w.Close()
|
|
|
|
first := time.Now().Add(80 * time.Millisecond) // would fire warning ~30ms in
|
|
_ = w.Update(first)
|
|
|
|
// Replace with a far-future deadline before the warning fires.
|
|
time.Sleep(5 * time.Millisecond)
|
|
second := time.Now().Add(time.Hour)
|
|
_ = w.Update(second)
|
|
|
|
// Wait past when first's warning would have fired.
|
|
time.Sleep(80 * time.Millisecond)
|
|
|
|
if n := countWhere(r.snapshot(), event.isWarning); n != 0 {
|
|
t.Fatalf("warning fired for cancelled deadline: %+v", r.snapshot())
|
|
}
|
|
}
|
|
|
|
func TestRefreshAfterFireArmsNewWarning(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
lead := 150 * time.Millisecond
|
|
w := newWatcher(lead, r)
|
|
defer w.Close()
|
|
|
|
// Warning fires ~20ms in; the deadline itself stays 150ms away so the
|
|
// replacement below lands well before it.
|
|
first := time.Now().Add(170 * time.Millisecond)
|
|
_ = w.Update(first)
|
|
|
|
// Wait for stateChange + warning of the first cycle.
|
|
waitForEvents(t, r, 2)
|
|
|
|
// Simulate a successful extend: brand new deadline.
|
|
second := time.Now().Add(60 * time.Millisecond)
|
|
_ = w.Update(second)
|
|
|
|
// 4 events total: stateChange, warning (first), stateChange, warning (second).
|
|
events := waitForEvents(t, r, 4)
|
|
if events[2].kind != stateChange {
|
|
t.Fatalf("event[2] should be stateChange for the new deadline, got %+v", events[2])
|
|
}
|
|
if !events[3].isWarning() {
|
|
t.Fatalf("event[3] should be a warning publish for the new deadline, got %+v", events[3])
|
|
}
|
|
}
|
|
|
|
func TestUpdateZeroAfterNonZeroClearsState(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(time.Hour, r)
|
|
defer w.Close()
|
|
|
|
d := time.Now().Add(2 * time.Hour)
|
|
_ = w.Update(d)
|
|
waitForEvents(t, r, 1)
|
|
|
|
_ = w.Update(time.Time{})
|
|
|
|
events := waitForEvents(t, r, 2)
|
|
if events[1].kind != stateChange {
|
|
t.Fatalf("expected stateChange on clear, got %+v", events[1])
|
|
}
|
|
if !w.Deadline().IsZero() {
|
|
t.Fatalf("Deadline should be zero after clear")
|
|
}
|
|
}
|
|
|
|
func TestUpdateRejectsBeforeEpoch(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(50*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
good := time.Now().Add(time.Hour)
|
|
if err := w.Update(good); err != nil {
|
|
t.Fatalf("seed Update: %v", err)
|
|
}
|
|
|
|
err := w.Update(time.Unix(-100, 0))
|
|
if !errors.Is(err, ErrDeadlineBeforeEpoch) {
|
|
t.Fatalf("want ErrDeadlineBeforeEpoch, got %v", err)
|
|
}
|
|
if !w.Deadline().IsZero() {
|
|
t.Fatalf("rejected pre-epoch update must clear deadline; got %v", w.Deadline())
|
|
}
|
|
}
|
|
|
|
func TestUpdateRejectsTooFarFuture(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(50*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
good := time.Now().Add(time.Hour)
|
|
if err := w.Update(good); err != nil {
|
|
t.Fatalf("seed Update: %v", err)
|
|
}
|
|
|
|
err := w.Update(time.Now().Add(50 * 365 * 24 * time.Hour))
|
|
if !errors.Is(err, ErrDeadlineTooFarFuture) {
|
|
t.Fatalf("want ErrDeadlineTooFarFuture, got %v", err)
|
|
}
|
|
if !w.Deadline().IsZero() {
|
|
t.Fatalf("rejected far-future update must clear deadline; got %v", w.Deadline())
|
|
}
|
|
}
|
|
|
|
func TestUpdateRecentPastRecordedAsExpired(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(50*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
d := time.Now().Add(-1 * time.Hour)
|
|
if err := w.Update(d); err != nil {
|
|
t.Fatalf("recent-past Update should succeed, got %v", err)
|
|
}
|
|
if !w.Deadline().Equal(d) {
|
|
t.Fatalf("expected deadline to be recorded, got %v want %v", w.Deadline(), d)
|
|
}
|
|
if got := r.deadline(); !got.Equal(d) {
|
|
t.Fatalf("recorder deadline = %v, want %v", got, d)
|
|
}
|
|
|
|
time.Sleep(80 * time.Millisecond)
|
|
if n := countWhere(r.snapshot(), func(e event) bool { return e.kind == publish }); n != 0 {
|
|
t.Fatalf("no warning events may fire for an already-past deadline, got %+v", r.snapshot())
|
|
}
|
|
}
|
|
|
|
func TestUpdateAncientPastRejected(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(50*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
good := time.Now().Add(time.Hour)
|
|
if err := w.Update(good); err != nil {
|
|
t.Fatalf("seed Update: %v", err)
|
|
}
|
|
// Drain the stateChange from the seed.
|
|
waitForEvents(t, r, 1)
|
|
|
|
err := w.Update(time.Now().Add(-31 * 24 * time.Hour))
|
|
if !errors.Is(err, ErrDeadlineInPast) {
|
|
t.Fatalf("want ErrDeadlineInPast, got %v", err)
|
|
}
|
|
if !w.Deadline().IsZero() {
|
|
t.Fatalf("rejected ancient-past update must clear the deadline, got %v", w.Deadline())
|
|
}
|
|
events := waitForEvents(t, r, 2)
|
|
if events[1].kind != stateChange {
|
|
t.Fatalf("expected stateChange on clear, got %+v", events[1])
|
|
}
|
|
}
|
|
|
|
func TestCloseSilencesUpdates(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(50*time.Millisecond, r)
|
|
w.Close()
|
|
|
|
if err := w.Update(time.Now().Add(time.Hour)); err != nil {
|
|
t.Fatalf("Update after Close: want nil, got %v", err)
|
|
}
|
|
if got := r.snapshot(); len(got) != 0 {
|
|
t.Fatalf("expected no events after Close, got %+v", got)
|
|
}
|
|
}
|
|
|
|
// TestCloseKeepsRecorderDeadline pins the reconnect-flap fix: the watcher
|
|
// closes on every engine restart (network change, sleep/wake) while the
|
|
// SSO deadline stays valid across those, so Close must leave the
|
|
// server-scoped recorder's value in place. The client run loop clears the
|
|
// recorder when it exits for real.
|
|
func TestCloseKeepsRecorderDeadline(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(time.Hour, r)
|
|
|
|
d := time.Now().Add(2 * time.Hour)
|
|
if err := w.Update(d); err != nil {
|
|
t.Fatalf("seed Update: %v", err)
|
|
}
|
|
if got := r.deadline(); !got.Equal(d) {
|
|
t.Fatalf("recorder deadline after Update = %v, want %v", got, d)
|
|
}
|
|
|
|
w.Close()
|
|
|
|
if got := r.deadline(); !got.Equal(d) {
|
|
t.Fatalf("recorder deadline after Close = %v, want %v", got, d)
|
|
}
|
|
}
|
|
|
|
// TestCloseWithoutDeadlineLeavesRecorderUntouched guards the symmetric
|
|
// case: closing a watcher that never held a deadline must not emit a
|
|
// redundant clear (the recorder may legitimately hold a value written by
|
|
// some other path; the watcher only owns what it set).
|
|
func TestCloseWithoutDeadlineLeavesRecorderUntouched(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcher(time.Hour, r)
|
|
|
|
w.Close()
|
|
|
|
if got := r.snapshot(); len(got) != 0 {
|
|
t.Fatalf("expected no events from Close on an empty watcher, got %+v", got)
|
|
}
|
|
}
|
|
|
|
func TestFinalWarningFiresAfterRegularWarning(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
// Warning fires at deadline-80ms, final at deadline-30ms.
|
|
w := newWatcherWithLeads(80*time.Millisecond, 30*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
d := time.Now().Add(100 * time.Millisecond)
|
|
_ = w.Update(d)
|
|
|
|
// Expect stateChange + warning + final-warning.
|
|
events := waitForEvents(t, r, 3)
|
|
|
|
if countWhere(events, func(e event) bool { return e.kind == stateChange }) != 1 {
|
|
t.Fatalf("expected exactly 1 stateChange, got %+v", events)
|
|
}
|
|
if countWhere(events, event.isWarning) != 1 {
|
|
t.Fatalf("expected exactly 1 warning publish, got %+v", events)
|
|
}
|
|
if countWhere(events, event.isFinalWarning) != 1 {
|
|
t.Fatalf("expected exactly 1 final-warning publish, got %+v", events)
|
|
}
|
|
|
|
// Warning must precede final (same deadline, longer lead fires first).
|
|
var wIdx, fIdx int
|
|
for i, e := range events {
|
|
switch {
|
|
case e.isWarning():
|
|
wIdx = i
|
|
case e.isFinalWarning():
|
|
fIdx = i
|
|
}
|
|
}
|
|
if wIdx > fIdx {
|
|
t.Fatalf("warning must publish before final-warning, got order %+v", events)
|
|
}
|
|
}
|
|
|
|
func TestDismissSuppressesFinalWarning(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcherWithLeads(80*time.Millisecond, 30*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
d := time.Now().Add(100 * time.Millisecond)
|
|
_ = w.Update(d)
|
|
|
|
// Wait for the warning publish so we know we're inside the warning
|
|
// window, then dismiss before the final timer would fire.
|
|
deadline := time.Now().Add(500 * time.Millisecond)
|
|
for time.Now().Before(deadline) {
|
|
if countWhere(r.snapshot(), event.isWarning) >= 1 {
|
|
break
|
|
}
|
|
time.Sleep(2 * time.Millisecond)
|
|
}
|
|
if countWhere(r.snapshot(), event.isWarning) < 1 {
|
|
t.Fatalf("warning did not publish in time, events=%+v", r.snapshot())
|
|
}
|
|
|
|
w.Dismiss()
|
|
|
|
// Now wait past when the final would have fired.
|
|
time.Sleep(120 * time.Millisecond)
|
|
|
|
if n := countWhere(r.snapshot(), event.isFinalWarning); n != 0 {
|
|
t.Fatalf("final-warning published after Dismiss(), events=%+v", r.snapshot())
|
|
}
|
|
}
|
|
|
|
func TestDismissResetByNewDeadline(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcherWithLeads(80*time.Millisecond, 30*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
first := time.Now().Add(100 * time.Millisecond)
|
|
_ = w.Update(first)
|
|
|
|
// Dismiss against the first deadline.
|
|
w.Dismiss()
|
|
|
|
// Replace with a fresh deadline before the first's timers complete.
|
|
time.Sleep(10 * time.Millisecond)
|
|
second := time.Now().Add(100 * time.Millisecond)
|
|
_ = w.Update(second)
|
|
|
|
// The second cycle must publish a final-warning (the dismiss state
|
|
// did not carry over).
|
|
deadline := time.Now().Add(500 * time.Millisecond)
|
|
for time.Now().Before(deadline) {
|
|
if countWhere(r.snapshot(), event.isFinalWarning) >= 1 {
|
|
break
|
|
}
|
|
time.Sleep(5 * time.Millisecond)
|
|
}
|
|
if countWhere(r.snapshot(), event.isFinalWarning) < 1 {
|
|
t.Fatalf("final-warning did not publish on fresh deadline after Dismiss reset, events=%+v", r.snapshot())
|
|
}
|
|
}
|
|
|
|
func TestDismissBeforeUpdateIsNoop(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcherWithLeads(80*time.Millisecond, 30*time.Millisecond, r)
|
|
defer w.Close()
|
|
|
|
// No deadline tracked yet; Dismiss must be a no-op (no panic, no state).
|
|
w.Dismiss()
|
|
|
|
d := time.Now().Add(100 * time.Millisecond)
|
|
_ = w.Update(d)
|
|
|
|
// Final warning should still publish — Dismiss only acts on the current
|
|
// deadline, and there was none at the time of the call.
|
|
deadline := time.Now().Add(500 * time.Millisecond)
|
|
for time.Now().Before(deadline) {
|
|
if countWhere(r.snapshot(), event.isFinalWarning) >= 1 {
|
|
return
|
|
}
|
|
time.Sleep(5 * time.Millisecond)
|
|
}
|
|
t.Fatalf("final-warning did not publish after no-op pre-Update Dismiss, events=%+v", r.snapshot())
|
|
}
|
|
|
|
// The tests below drive the watcher's wall clock directly. The deadline sits
|
|
// an hour out in real time and the fake clock jumps, standing in for a device
|
|
// that was suspended and resumed somewhere inside — or past — the warning
|
|
// windows. The evaluation loop is what reacts to the jump, so these exercise
|
|
// the same path production takes on a resume.
|
|
func newResumeWatcher(t *testing.T, r *fakeRecorder) (*Watcher, *fakeClock, time.Time) {
|
|
t.Helper()
|
|
|
|
w := newWatcherWithLeads(WarningLead, FinalWarningLead, r)
|
|
start := time.Now()
|
|
clock := newFakeClock(start)
|
|
// Set before Update: the evaluation loop does not exist yet.
|
|
w.nowFn = clock.now
|
|
t.Cleanup(w.Close)
|
|
|
|
deadline := start.Add(time.Hour).Round(0)
|
|
if err := w.Update(deadline); err != nil {
|
|
t.Fatalf("Update: %v", err)
|
|
}
|
|
return w, clock, deadline
|
|
}
|
|
|
|
func TestResumeInsideWarningWindowWarns(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
_, clock, deadline := newResumeWatcher(t, r)
|
|
|
|
clock.set(deadline.Add(-5 * time.Minute))
|
|
|
|
events := waitForEvents(t, r, 2)
|
|
if !events[1].isWarning() {
|
|
t.Fatalf("expected the interactive warning after the resume, got %+v", events[1])
|
|
}
|
|
if n := countWhere(events, event.isFinalWarning); n != 0 {
|
|
t.Fatalf("final-warning must wait for its own window, got %d: %+v", n, events)
|
|
}
|
|
}
|
|
|
|
func TestResumeInsideFinalWindowSendsFinalWarningOnly(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
_, clock, deadline := newResumeWatcher(t, r)
|
|
|
|
clock.set(deadline.Add(-time.Minute))
|
|
|
|
events := waitForEvents(t, r, 2)
|
|
if !events[1].isFinalWarning() {
|
|
t.Fatalf("expected the final warning after the resume, got %+v", events[1])
|
|
}
|
|
settle()
|
|
if n := countWhere(r.snapshot(), event.isWarning); n != 0 {
|
|
t.Fatalf("the interactive warning is stale inside the final window, got %d: %+v", n, r.snapshot())
|
|
}
|
|
}
|
|
|
|
func TestResumePastDeadlinePublishesNothing(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
_, clock, deadline := newResumeWatcher(t, r)
|
|
|
|
clock.set(deadline.Add(time.Minute))
|
|
|
|
settle()
|
|
if n := countWhere(r.snapshot(), func(e event) bool { return e.kind == publish }); n != 0 {
|
|
t.Fatalf("an expired session must not warn, got %d publishes: %+v", n, r.snapshot())
|
|
}
|
|
}
|
|
|
|
func TestWarningPublishesOncePerDeadline(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
_, clock, deadline := newResumeWatcher(t, r)
|
|
|
|
clock.set(deadline.Add(-5 * time.Minute))
|
|
waitForEvents(t, r, 2)
|
|
|
|
// Many ticks pass inside the same window.
|
|
settle()
|
|
if n := countWhere(r.snapshot(), event.isWarning); n != 1 {
|
|
t.Fatalf("expected exactly 1 warning publish across ticks, got %d: %+v", n, r.snapshot())
|
|
}
|
|
}
|
|
|
|
// TestWarningRecoversFromAClockRunningAhead covers the device that boots
|
|
// before NTP has corrected it: the deadline looks long gone, nothing is
|
|
// published, and the warning still arrives once the clock is fixed.
|
|
func TestWarningRecoversFromAClockRunningAhead(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
_, clock, deadline := newResumeWatcher(t, r)
|
|
|
|
clock.set(deadline.Add(2 * time.Hour))
|
|
settle()
|
|
if n := countWhere(r.snapshot(), func(e event) bool { return e.kind == publish }); n != 0 {
|
|
t.Fatalf("a deadline that looks expired must not warn, got %d publishes: %+v", n, r.snapshot())
|
|
}
|
|
|
|
clock.set(deadline.Add(-5 * time.Minute))
|
|
|
|
events := waitForEvents(t, r, 2)
|
|
if !events[1].isWarning() {
|
|
t.Fatalf("expected the warning once the clock was corrected, got %+v", events[1])
|
|
}
|
|
}
|
|
|
|
func TestResumeInsideFinalWindowRespectsDismiss(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w, clock, deadline := newResumeWatcher(t, r)
|
|
|
|
// The user dismissed the warning for this deadline before the device
|
|
// was suspended; resuming inside the final window must not reopen it.
|
|
w.Dismiss()
|
|
clock.set(deadline.Add(-time.Minute))
|
|
|
|
settle()
|
|
if n := countWhere(r.snapshot(), func(e event) bool { return e.kind == publish }); n != 0 {
|
|
t.Fatalf("a dismissed deadline must not warn on resume, got %d: %+v", n, r.snapshot())
|
|
}
|
|
}
|
|
|
|
func TestCloseStopsTheEvaluationLoop(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcherWithLeads(WarningLead, FinalWarningLead, r)
|
|
start := time.Now()
|
|
clock := newFakeClock(start)
|
|
w.nowFn = clock.now
|
|
|
|
deadline := start.Add(time.Hour).Round(0)
|
|
if err := w.Update(deadline); err != nil {
|
|
t.Fatalf("Update: %v", err)
|
|
}
|
|
w.Close()
|
|
|
|
clock.set(deadline.Add(-5 * time.Minute))
|
|
settle()
|
|
|
|
if n := countWhere(r.snapshot(), func(e event) bool { return e.kind == publish }); n != 0 {
|
|
t.Fatalf("a closed watcher must not publish, got %d: %+v", n, r.snapshot())
|
|
}
|
|
}
|
|
|
|
func TestDeadlineOnlyRecordsDeadlineWithoutWarnings(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := NewDeadlineOnly(r)
|
|
defer w.Close()
|
|
|
|
// With the default leads this deadline sits inside the final-warning
|
|
// window, so a watcher that warns at all would publish on the spot.
|
|
d := time.Now().Add(50 * time.Millisecond).Round(0)
|
|
if err := w.Update(d); err != nil {
|
|
t.Fatalf("Update: %v", err)
|
|
}
|
|
if got := r.deadline(); !got.Equal(d) {
|
|
t.Fatalf("expected recorder deadline %v, got %v", d, got)
|
|
}
|
|
|
|
time.Sleep(100 * time.Millisecond)
|
|
|
|
events := r.snapshot()
|
|
if got := countWhere(events, func(e event) bool { return e.kind == publish }); got != 0 {
|
|
t.Fatalf("expected no publish in deadline-only mode, got %d: %+v", got, events)
|
|
}
|
|
w.mu.Lock()
|
|
polling := w.stop != nil
|
|
w.mu.Unlock()
|
|
if polling {
|
|
t.Fatal("expected no evaluation loop in deadline-only mode")
|
|
}
|
|
}
|
|
|
|
func TestDeadlineOnlyStillRejectsOutOfRangeDeadlines(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := NewDeadlineOnly(r)
|
|
defer w.Close()
|
|
|
|
if err := w.Update(time.Now().Add(time.Hour)); err != nil {
|
|
t.Fatalf("Update: %v", err)
|
|
}
|
|
|
|
err := w.Update(time.Now().Add(-maxPastHorizon - time.Hour))
|
|
if !errors.Is(err, ErrDeadlineInPast) {
|
|
t.Fatalf("expected ErrDeadlineInPast, got %v", err)
|
|
}
|
|
if got := r.deadline(); !got.IsZero() {
|
|
t.Fatalf("expected recorder cleared after rejection, got %v", got)
|
|
}
|
|
}
|
|
|
|
// TestClockAheadAtUpdateStillWarnsOnceCorrected covers a client that connects
|
|
// while its clock runs ahead of real time: the deadline reads as already
|
|
// expired when it arrives, so nothing publishes, and the warning must still
|
|
// come once NTP pulls the clock back.
|
|
func TestClockAheadAtUpdateStillWarnsOnceCorrected(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcherWithLeads(WarningLead, FinalWarningLead, r)
|
|
deadline := time.Now().Add(time.Hour).Round(0)
|
|
// Two hours past the deadline from the device's point of view, well
|
|
// inside maxPastHorizon, so Update accepts and records it.
|
|
clock := newFakeClock(deadline.Add(2 * time.Hour))
|
|
w.nowFn = clock.now
|
|
t.Cleanup(w.Close)
|
|
|
|
if err := w.Update(deadline); err != nil {
|
|
t.Fatalf("Update: %v", err)
|
|
}
|
|
|
|
settle()
|
|
if n := countWhere(r.snapshot(), func(e event) bool { return e.kind == publish }); n != 0 {
|
|
t.Fatalf("a deadline that reads as expired must not warn, got %d: %+v", n, r.snapshot())
|
|
}
|
|
|
|
clock.set(deadline.Add(-5 * time.Minute))
|
|
|
|
events := waitForEvents(t, r, 2)
|
|
if !events[1].isWarning() {
|
|
t.Fatalf("expected the warning once the clock was corrected, got %+v", events[1])
|
|
}
|
|
}
|
|
|
|
// TestWarningNeverPrecedesTheDeadlineStateChange pins the ordering Update
|
|
// documents: consumers learn the new deadline before they see a warning that
|
|
// refers to it. The deadline here already sits inside the warning window, and
|
|
// the recorder stalls, so the evaluation loop gets many chances to publish
|
|
// while Update is still announcing.
|
|
func TestWarningNeverPrecedesTheDeadlineStateChange(t *testing.T) {
|
|
r := &fakeRecorder{setDelay: 50 * time.Millisecond}
|
|
w := newWatcherWithLeads(WarningLead, FinalWarningLead, r)
|
|
deadline := time.Now().Add(time.Hour).Round(0)
|
|
clock := newFakeClock(deadline.Add(-5 * time.Minute))
|
|
w.nowFn = clock.now
|
|
t.Cleanup(w.Close)
|
|
|
|
if err := w.Update(deadline); err != nil {
|
|
t.Fatalf("Update: %v", err)
|
|
}
|
|
|
|
events := waitForEvents(t, r, 2)
|
|
if events[0].kind != stateChange {
|
|
t.Fatalf("event[0] should be the deadline state change, got %+v", events)
|
|
}
|
|
if !events[1].isWarning() {
|
|
t.Fatalf("event[1] should be the warning, got %+v", events)
|
|
}
|
|
}
|
|
|
|
// TestUpdateWakesTheLoopWithoutWaitingForATick pins the hand-off: Update
|
|
// publishes nothing itself, it nudges the evaluation loop, so a deadline that
|
|
// already sits inside a warning window is warned about at once even though the
|
|
// ticker here would not fire for an hour.
|
|
func TestUpdateWakesTheLoopWithoutWaitingForATick(t *testing.T) {
|
|
r := &fakeRecorder{}
|
|
w := newWatcherWithLeads(WarningLead, FinalWarningLead, r)
|
|
w.interval = time.Hour
|
|
t.Cleanup(w.Close)
|
|
|
|
d := time.Now().Add(5 * time.Minute).Round(0)
|
|
if err := w.Update(d); err != nil {
|
|
t.Fatalf("Update: %v", err)
|
|
}
|
|
|
|
events := waitForEvents(t, r, 2)
|
|
if !events[1].isWarning() {
|
|
t.Fatalf("expected the warning on the wake-up, got %+v", events[1])
|
|
}
|
|
}
|
|
|
|
// blockingRecorder holds a publish open until the test releases it, so the
|
|
// evaluation loop can be parked mid-publish while Close is called.
|
|
type blockingRecorder struct {
|
|
fakeRecorder
|
|
entered chan struct{}
|
|
release chan struct{}
|
|
}
|
|
|
|
func (r *blockingRecorder) PublishEvent(
|
|
severity cProto.SystemEvent_Severity,
|
|
category cProto.SystemEvent_Category,
|
|
message string,
|
|
userMessage string,
|
|
metadata map[string]string,
|
|
) {
|
|
select {
|
|
case r.entered <- struct{}{}:
|
|
default:
|
|
}
|
|
<-r.release
|
|
r.fakeRecorder.PublishEvent(severity, category, message, userMessage, metadata)
|
|
}
|
|
|
|
// TestConcurrentCloseWaitsForTheLoop pins the contract for the caller that
|
|
// loses the race: Close returns only once the loop is done, whichever of the
|
|
// two calls got there first.
|
|
func TestConcurrentCloseWaitsForTheLoop(t *testing.T) {
|
|
r := &blockingRecorder{entered: make(chan struct{}, 1), release: make(chan struct{})}
|
|
var releaseOnce sync.Once
|
|
release := func() { releaseOnce.Do(func() { close(r.release) }) }
|
|
|
|
w := newWatcherWithLeads(WarningLead, FinalWarningLead, r)
|
|
w.interval = time.Hour
|
|
// Order matters: Close waits for the loop, which stays parked in the
|
|
// publish until the release, so a t.Fatal below would hang the cleanup
|
|
// instead of reporting the failure.
|
|
t.Cleanup(func() {
|
|
release()
|
|
w.Close()
|
|
})
|
|
|
|
d := time.Now().Add(5 * time.Minute).Round(0)
|
|
if err := w.Update(d); err != nil {
|
|
t.Fatalf("Update: %v", err)
|
|
}
|
|
|
|
// The loop is now parked inside the warning publish.
|
|
select {
|
|
case <-r.entered:
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("the loop never reached the publish")
|
|
}
|
|
|
|
first := make(chan struct{})
|
|
second := make(chan struct{})
|
|
go func() { w.Close(); close(first) }()
|
|
go func() { w.Close(); close(second) }()
|
|
|
|
select {
|
|
case <-first:
|
|
t.Fatal("Close returned while the loop was still publishing")
|
|
case <-second:
|
|
t.Fatal("Close returned while the loop was still publishing")
|
|
case <-time.After(100 * time.Millisecond):
|
|
}
|
|
|
|
release()
|
|
for _, done := range []chan struct{}{first, second} {
|
|
select {
|
|
case <-done:
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("Close did not return after the publish completed")
|
|
}
|
|
}
|
|
}
|