diff --git a/client/internal/auth/sessionwatch/watcher.go b/client/internal/auth/sessionwatch/watcher.go index 496903044..f71689df1 100644 --- a/client/internal/auth/sessionwatch/watcher.go +++ b/client/internal/auth/sessionwatch/watcher.go @@ -4,12 +4,20 @@ // T-WarningLead notification and a dismiss-gated T-FinalWarningLead // fallback dialog. // +// The deadline is an absolute wall-clock instant, so the watcher compares +// it against the wall clock on a ticker rather than arming a relative +// timer for it. 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. Polling makes every tick after a +// resume (or after an NTP correction) see the real remaining time. +// // The watcher is idempotent: Update may be called as often as the network // map snapshots arrive. Repeating the same deadline is a no-op; a new -// deadline reschedules the timers and arms a fresh warning cycle. +// deadline starts a fresh warning cycle. // -// Warning firing is edge-detected. Each unique deadline value fires each -// warning callback at most once. +// Warning firing is edge-detected. Each unique deadline value publishes +// each warning at most once. package sessionwatch import ( @@ -28,9 +36,15 @@ const ( // maxDeadlineHorizon caps how far in the future an accepted deadline // can sit. A timestamp beyond this is almost certainly a protocol - // glitch, and silently arming a 100-year timer would hide the bug. + // glitch, and silently tracking a 100-year deadline would hide the bug. maxDeadlineHorizon = 10 * 365 * 24 * time.Hour + // defaultEvalInterval is how often the tracked deadline is compared + // against the wall clock. The leads are minutes, so a coarse tick costs + // nothing in accuracy and keeps the wakeup cheap on battery-powered + // devices. + defaultEvalInterval = 10 * time.Second + // WarningLead is how far before expiry the first (interactive) // warning fires. Drives the T-10 OS notification with // Extend/Dismiss actions. @@ -92,18 +106,21 @@ type StatusRecorder interface { type Watcher struct { lead time.Duration finalLead time.Duration - deadlineOnly bool + interval time.Duration + deadlineOnly bool // record the deadline, leave the warnings to the caller mu sync.Mutex current time.Time - timer *time.Timer - finalTimer *time.Timer - firedAt time.Time // deadline value the T-WarningLead callback last fired against - finalFiredAt time.Time // deadline value the T-FinalWarningLead callback last fired against - dismissedAt time.Time // deadline value the user dismissed via Dismiss(); gates fireFinal + firedAt time.Time // deadline value the T-WarningLead warning last published for + finalFiredAt time.Time // deadline value the T-FinalWarningLead warning last published for + dismissedAt time.Time // deadline value the user dismissed via Dismiss(); gates the final warning + announcedAt time.Time // deadline value the recorder has been told about; gates publishing closed bool recorder StatusRecorder nowFn func() time.Time + stop chan struct{} // closed to stop the evaluation loop; nil while it is not running + done chan struct{} // closed by the loop on its way out + wake chan struct{} // buffered nudge asking the loop to evaluate before its next tick } // New returns a watcher with the package defaults WarningLead and @@ -115,20 +132,24 @@ func New(recorder StatusRecorder) *Watcher { } // NewWithLeads returns a watcher with custom lead times. Useful for tests. -// final must be strictly less than lead; otherwise both timers fire in the -// wrong order or simultaneously and the UI flow breaks. A zero final lead -// disables the final-warning timer entirely (see armTimerLocked) so a -// millisecond-scale deadline doesn't flush both timers in one tick. +// final must be strictly less than lead; otherwise the final warning takes +// over the whole warning window and the interactive notification never +// shows. A zero final lead disables the final warning entirely (see +// evaluate), leaving the interactive one as the only warning. func NewWithLeads(lead, final time.Duration, recorder StatusRecorder) *Watcher { return &Watcher{ lead: lead, finalLead: final, + interval: defaultEvalInterval, recorder: recorder, nowFn: time.Now, } } -// NewDeadlineOnly returns a watcher that validates and records deadlines but arms no warning timers. +// NewDeadlineOnly returns a watcher that validates and records deadlines +// but publishes no warnings about them, and runs no evaluation loop to +// decide. Used where the deadline is handed on to something that schedules +// the warnings itself, such as the Android app. func NewDeadlineOnly(recorder StatusRecorder) *Watcher { w := New(recorder) w.deadlineOnly = true @@ -139,11 +160,12 @@ func NewDeadlineOnly(recorder StatusRecorder) *Watcher { // a Sync push from the server omits the field because login expiration // was disabled). // -// Same-value updates are no-ops. A different non-zero value cancels any -// pending timer, resets the "already fired" guards, and — when the -// deadline lies in the future — arms fresh warning timers. A deadline -// already in the past (within maxPastHorizon) is recorded as-is with no -// timers: the session has expired and consumers render it that way. +// Same-value updates are no-ops. A different non-zero value resets the +// "already fired" guards and starts a fresh warning cycle, evaluated +// immediately so a deadline that already sits inside a warning window +// warns without waiting for the next tick. A deadline already in the past +// (within maxPastHorizon) is recorded as-is and warns nothing: the session +// has expired and consumers render it that way. // // Returns one of the sentinel Err* values when the deadline fails the // sanity checks (pre-epoch, far future, or past beyond maxPastHorizon). @@ -163,7 +185,7 @@ func (w *Watcher) Update(deadline time.Time) error { return nil } - now := time.Now() + now := w.nowFn() switch { case deadline.Before(time.Unix(0, 0)): w.clearLocked() @@ -181,18 +203,22 @@ func (w *Watcher) Update(deadline time.Time) error { return nil } - w.stopTimerLocked() w.current = deadline - // Reset every per-deadline guard so a refreshed deadline arms a fresh + // Reset every per-deadline guard so a refreshed deadline starts a fresh // warning cycle: both edge triggers and the user Dismiss decision // (the user agreed to the old deadline expiring; a new deadline // restarts the contract). w.firedAt = time.Time{} w.finalFiredAt = time.Time{} w.dismissedAt = time.Time{} + w.announcedAt = time.Time{} - if deadline.After(now) && !w.deadlineOnly { - w.armTimerLocked(deadline) + // Poll every accepted deadline, including one that reads as already + // expired: the clock may be running ahead of real time and get + // corrected later, and evaluate ignores a deadline that has genuinely + // passed anyway. + if !w.deadlineOnly { + w.startPollLocked() } recorder := w.recorder w.mu.Unlock() @@ -200,6 +226,28 @@ func (w *Watcher) Update(deadline time.Time) error { recorder.SetSessionExpiresAt(deadline) } log.Infof("auth session deadline set to: %s (in %s)", deadline.Format(time.RFC3339), time.Until(deadline).Round(time.Second)) + + // Open the gate only once the recorder knows the new deadline, so a + // warning that refers to it can never reach consumers before the state + // change itself. A tick landing in between finds the gate shut. + w.mu.Lock() + if w.closed || !w.current.Equal(deadline) { + w.mu.Unlock() + return nil + } + w.announcedAt = deadline + wake := w.wake + w.mu.Unlock() + + // Hand the evaluation to the loop rather than running it here, so a + // deadline that already sits inside a warning window is published at + // once without this goroutine ever touching the recorder. + if wake != nil { + select { + case wake <- struct{}{}: + default: + } + } return nil } @@ -212,9 +260,9 @@ func (w *Watcher) Deadline() time.Time { } // Dismiss records the user's "Dismiss" action against the current deadline -// and suppresses the upcoming final-warning callback for that deadline. -// Idempotent: repeated calls are no-ops. A subsequent Update with a fresh -// deadline resets the dismissal so the final-warning cycle re-arms. +// and suppresses the final warning for that deadline. Idempotent: repeated +// calls are no-ops. A subsequent Update with a fresh deadline resets the +// dismissal so the final-warning cycle starts over. // // No-op when the watcher holds no deadline or has been closed. func (w *Watcher) Dismiss() { @@ -227,35 +275,48 @@ func (w *Watcher) Dismiss() { return } w.dismissedAt = w.current - // Cancel the armed final-warning timer eagerly. fireFinal would also - // gate on dismissedAt, but stopping the timer avoids a wakeup with - // nothing to do and makes the intent visible. - if w.finalTimer != nil { - w.finalTimer.Stop() - w.finalTimer = nil - } log.Infof("auth session final-warning dismissed for deadline %s", w.current.Format(time.RFC3339)) } -// Close stops any pending timer. Update calls after Close are ignored. -// The recorder keeps its deadline: the watcher is engine-scoped and closes -// on every engine restart (network change, sleep/wake, stream errors) -// while the SSO deadline stays valid across those, so clearing here would -// blank the UI's "expires in" row on every transient reconnect. The -// client run loop clears the server-scoped recorder when it exits for -// real (Down, profile switch, permanent login failure). +// Close stops the evaluation loop and waits for it to exit. Update calls +// after Close are ignored. The recorder keeps its deadline: the watcher is +// engine-scoped and closes on every engine restart (network change, +// sleep/wake, stream errors) while the SSO deadline stays valid across +// those, so clearing here would blank the UI's "expires in" row on every +// transient reconnect. The client run loop clears the server-scoped +// recorder when it exits for real (Down, profile switch, permanent login +// failure). func (w *Watcher) Close() { w.mu.Lock() - defer w.mu.Unlock() if w.closed { + // A concurrent Close is already tearing down. w.done outlives it, + // so this caller waits for the same loop rather than returning + // while a warning is still on its way to the recorder. + done := w.done + w.mu.Unlock() + if done != nil { + <-done + } return } w.closed = true - w.stopTimerLocked() w.current = time.Time{} w.firedAt = time.Time{} w.finalFiredAt = time.Time{} w.dismissedAt = time.Time{} + w.announcedAt = time.Time{} + // Copy the channels out before releasing the lock: the loop takes w.mu + // on every tick, so waiting for it while holding the lock would + // deadlock. w.done stays on the receiver for the branch above. + stop, done := w.stop, w.done + w.stop, w.wake = nil, nil + w.mu.Unlock() + + if stop == nil { + return + } + close(stop) + <-done } // clearLocked drops the tracked deadline and notifies the recorder so @@ -267,11 +328,11 @@ func (w *Watcher) clearLocked() { w.mu.Unlock() return } - w.stopTimerLocked() w.current = time.Time{} w.firedAt = time.Time{} w.finalFiredAt = time.Time{} w.dismissedAt = time.Time{} + w.announcedAt = time.Time{} recorder := w.recorder w.mu.Unlock() if recorder != nil { @@ -280,133 +341,121 @@ func (w *Watcher) clearLocked() { log.Infof("auth session deadline cleared") } -func (w *Watcher) stopTimerLocked() { - if w.timer != nil { - w.timer.Stop() - w.timer = nil +// startPollLocked starts the evaluation loop unless it is already +// running. The loop starts lazily on the first accepted deadline, so a +// client whose server never publishes a session expiry never pays for a +// ticker, and it runs until Close: the watcher is engine-scoped, and a +// cleared deadline is normally followed by a fresh one on the next sync. +// Caller must hold w.mu. +func (w *Watcher) startPollLocked() { + if w.stop != nil { + return } - if w.finalTimer != nil { - w.finalTimer.Stop() - w.finalTimer = nil + stop := make(chan struct{}) + done := make(chan struct{}) + wake := make(chan struct{}, 1) + w.stop, w.done, w.wake = stop, done, wake + go w.poll(stop, done, wake, w.interval) +} + +// poll re-evaluates the tracked deadline every interval, and as soon as a +// new deadline is announced. It is the only caller of evaluate, so Close +// waiting for it to exit is enough to know no warning is still on its way +// to the recorder. Its channels and interval are passed in rather than +// read off the receiver, so Close can clear them without racing it. +func (w *Watcher) poll(stop <-chan struct{}, done chan<- struct{}, wake <-chan struct{}, interval time.Duration) { + defer close(done) + + ticker := time.NewTicker(interval) + defer ticker.Stop() + + for { + select { + case <-stop: + return + case <-wake: + w.evaluate() + case <-ticker.C: + w.evaluate() + } } } -func (w *Watcher) armTimerLocked(deadline time.Time) { - w.timer = armOneShotLocked(deadline.Add(-w.lead), func() { w.fire(deadline) }) - // finalLead <= 0 disables the final-warning timer entirely. Used by - // tests that predate the final-warning fallback so a millisecond-scale - // deadline does not flush both timers at once. - if w.finalLead > 0 { - w.finalTimer = armOneShotLocked(deadline.Add(-w.finalLead), func() { w.fireFinal(deadline) }) - } -} - -func (w *Watcher) fire(armedFor time.Time) { +// evaluate compares the tracked deadline against the wall clock and +// publishes whichever warning the remaining time calls for. Inside the +// final-warning window the interactive warning is stale, so the final one +// is published in its place: a device that resumes there never saw the +// T-WarningLead notification. +func (w *Watcher) evaluate() { w.mu.Lock() - if w.closed || !w.current.Equal(armedFor) { - // Deadline moved while we were waiting (e.g. a successful extend). - // The reschedule path armed a fresh timer; this one is stale. + if w.closed || w.deadlineOnly || w.current.IsZero() { w.mu.Unlock() return } - if !w.firedAt.IsZero() && w.firedAt.Equal(armedFor) { - w.mu.Unlock() - return - } - now := w.nowFn() - if isLate(now, armedFor, max(w.finalLead, 0)) { - w.fireLateLocked(armedFor, now) - return - } - w.firedAt = armedFor - recorder := w.recorder - w.mu.Unlock() - if recorder == nil { - return - } - log.Infof("auth session expiry soon warning fired") - publishWarning(recorder, armedFor, false) -} -// fireFinal mirrors fire for the T-FinalWarningLead timer with an extra -// dismiss-gate: if the user dismissed the T-WarningLead notification for -// this deadline, the final warning is suppressed entirely. -func (w *Watcher) fireFinal(armedFor time.Time) { - w.mu.Lock() - if w.closed || !w.current.Equal(armedFor) { + deadline := w.current + if !w.announcedAt.Equal(deadline) { + // Update is still on its way to the recorder with this deadline. w.mu.Unlock() return } - if !w.finalFiredAt.IsZero() && w.finalFiredAt.Equal(armedFor) { - w.mu.Unlock() - return - } - if w.dismissedAt.Equal(armedFor) { - w.mu.Unlock() - log.Infof("auth session final-warning skipped (dismissed by user)") - return - } - now := w.nowFn() - if isLate(now, armedFor, 0) { - w.finalFiredAt = armedFor - w.mu.Unlock() - log.Infof("auth session final-warning skipped for deadline %s (passed %s ago)", - armedFor.Format(time.RFC3339), now.Round(0).Sub(armedFor).Round(time.Second)) - return - } - w.finalFiredAt = armedFor - recorder := w.recorder - w.mu.Unlock() - if recorder == nil { - return - } - log.Infof("auth session final-warning fired") - publishWarning(recorder, armedFor, true) -} + // Round(0) strips the monotonic reading so the comparison is wall + // clock on both sides, whether the deadline came off the wire or from + // a caller that derived it from time.Now. + remaining := deadline.Round(0).Sub(w.nowFn().Round(0)) -// fireLateLocked handles a T-WarningLead callback that fired inside the -// final-warning window: it sends the final warning in its place while the -// deadline has not passed and the user has not dismissed it, so a resume -// with time left still warns. The caller must hold w.mu; this helper -// releases it. -func (w *Watcher) fireLateLocked(armedFor, now time.Time) { - w.firedAt = armedFor switch { - case w.dismissedAt.Equal(armedFor): + case remaining <= 0: + // Already expired: the post-mortem SessionExpired flow owns it. w.mu.Unlock() - log.Infof("auth session expiry soon warning skipped (dismissed by user)") - return - case w.finalFiredAt.Equal(armedFor): + case w.finalLead > 0 && remaining <= w.finalLead: + w.publishFinalLocked(deadline, remaining) + case remaining <= w.lead: + w.publishWarningLocked(deadline, remaining) + default: w.mu.Unlock() - log.Infof("auth session expiry soon warning skipped (final warning already fired)") - return - case isLate(now, armedFor, 0): + } +} + +// publishWarningLocked emits the interactive T-WarningLead warning, at +// most once per deadline value. Caller must hold w.mu; this helper +// releases it. +func (w *Watcher) publishWarningLocked(deadline time.Time, remaining time.Duration) { + if w.firedAt.Equal(deadline) { w.mu.Unlock() - log.Infof("auth session expiry soon warning skipped for deadline %s (passed %s ago)", - armedFor.Format(time.RFC3339), now.Round(0).Sub(armedFor).Round(time.Second)) return } - w.finalFiredAt = armedFor + w.firedAt = deadline recorder := w.recorder w.mu.Unlock() if recorder == nil { return } - log.Infof("auth session expiry soon warning fired inside the final-warning window, sending final warning for deadline %s", - armedFor.Format(time.RFC3339)) - publishWarning(recorder, armedFor, true) + log.Infof("auth session expiry soon warning fired for deadline %s (in %s)", + deadline.Format(time.RFC3339), remaining.Round(time.Second)) + publishWarning(recorder, deadline, false) } -// armOneShotLocked schedules cb at fireAt. When fireAt is already in the -// past it dispatches on the next scheduler tick so a state-change recorder -// notification (invoked after w.mu is released) lands first. Caller must -// hold w.mu. -func armOneShotLocked(fireAt time.Time, cb func()) *time.Timer { - delay := time.Until(fireAt) - if delay <= 0 { - return time.AfterFunc(0, cb) +// publishFinalLocked emits the final warning, at most once per deadline +// value and never once the user dismissed that deadline. It marks the +// interactive warning as handled too: the final window is open, so a +// "expires in WarningLead minutes" notification would be wrong. Caller +// must hold w.mu; this helper releases it. +func (w *Watcher) publishFinalLocked(deadline time.Time, remaining time.Duration) { + if w.finalFiredAt.Equal(deadline) || w.dismissedAt.Equal(deadline) { + w.mu.Unlock() + return } - return time.AfterFunc(delay, cb) + w.firedAt = deadline + w.finalFiredAt = deadline + recorder := w.recorder + w.mu.Unlock() + if recorder == nil { + return + } + log.Infof("auth session final-warning fired for deadline %s (in %s)", + deadline.Format(time.RFC3339), remaining.Round(time.Second)) + publishWarning(recorder, deadline, true) } // publishWarning composes the SystemEvent for a watcher-fired warning and @@ -436,11 +485,3 @@ func publishWarning(recorder StatusRecorder, deadline time.Time, final bool) { meta, ) } - -// isLate reports whether the wall clock now has already reached armedFor -// minus cutoffLead. The timers run on the monotonic clock, which can stall -// while the host sleeps, so a timer can fire long after the window it was -// armed for. -func isLate(now, armedFor time.Time, cutoffLead time.Duration) bool { - return !now.Round(0).Before(armedFor.Add(-cutoffLead).Round(0)) -} diff --git a/client/internal/auth/sessionwatch/watcher_test.go b/client/internal/auth/sessionwatch/watcher_test.go index cb2800978..f0cc3ed88 100644 --- a/client/internal/auth/sessionwatch/watcher_test.go +++ b/client/internal/auth/sessionwatch/watcher_test.go @@ -19,6 +19,9 @@ 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 @@ -43,6 +46,12 @@ type event struct { // 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) { @@ -116,10 +125,54 @@ func waitForEvents(t *testing.T, r *fakeRecorder, want int) []event { return nil } -// newWatcher builds a watcher with the final timer disabled (finalLead=0), +// 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 NewWithLeads(lead, 0, r) + 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) { @@ -410,7 +463,7 @@ func TestCloseWithoutDeadlineLeavesRecorderUntouched(t *testing.T) { func TestFinalWarningFiresAfterRegularWarning(t *testing.T) { r := &fakeRecorder{} // Warning fires at deadline-80ms, final at deadline-30ms. - w := NewWithLeads(80*time.Millisecond, 30*time.Millisecond, r) + w := newWatcherWithLeads(80*time.Millisecond, 30*time.Millisecond, r) defer w.Close() d := time.Now().Add(100 * time.Millisecond) @@ -446,7 +499,7 @@ func TestFinalWarningFiresAfterRegularWarning(t *testing.T) { func TestDismissSuppressesFinalWarning(t *testing.T) { r := &fakeRecorder{} - w := NewWithLeads(80*time.Millisecond, 30*time.Millisecond, r) + w := newWatcherWithLeads(80*time.Millisecond, 30*time.Millisecond, r) defer w.Close() d := time.Now().Add(100 * time.Millisecond) @@ -477,7 +530,7 @@ func TestDismissSuppressesFinalWarning(t *testing.T) { func TestDismissResetByNewDeadline(t *testing.T) { r := &fakeRecorder{} - w := NewWithLeads(80*time.Millisecond, 30*time.Millisecond, r) + w := newWatcherWithLeads(80*time.Millisecond, 30*time.Millisecond, r) defer w.Close() first := time.Now().Add(100 * time.Millisecond) @@ -507,7 +560,7 @@ func TestDismissResetByNewDeadline(t *testing.T) { func TestDismissBeforeUpdateIsNoop(t *testing.T) { r := &fakeRecorder{} - w := NewWithLeads(80*time.Millisecond, 30*time.Millisecond, r) + 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). @@ -528,157 +581,139 @@ func TestDismissBeforeUpdateIsNoop(t *testing.T) { t.Fatalf("final-warning did not publish after no-op pre-Update Dismiss, events=%+v", r.snapshot()) } -func TestIsLate(t *testing.T) { - armedFor := time.Date(2026, 10, 1, 12, 0, 0, 0, time.UTC) - lead := 2 * time.Minute - tests := []struct { - name string - now time.Time - cutoffLead time.Duration - want bool - }{ - {"before cutoff", armedFor.Add(-3 * time.Minute), lead, false}, - {"at cutoff", armedFor.Add(-lead), lead, true}, - {"after cutoff", armedFor.Add(-time.Minute), lead, true}, - {"zero lead before deadline", armedFor.Add(-time.Second), 0, false}, - {"zero lead at deadline", armedFor, 0, true}, - {"zero lead after deadline", armedFor.Add(time.Second), 0, true}, - } - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - if got := isLate(tt.now, armedFor, tt.cutoffLead); got != tt.want { - t.Fatalf("isLate(%s, %s, %s) = %v, want %v", tt.now, armedFor, tt.cutoffLead, got, tt.want) - } - }) - } -} +// 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() -func TestIsLateIgnoresMonotonicReading(t *testing.T) { - now := time.Now() - wallOnly := now.Round(0) - if isLate(now, wallOnly.Add(time.Second), 0) { - t.Fatalf("now with monotonic reading must compare as wall clock before a later wall-only deadline") - } - if !isLate(now, wallOnly, 0) { - t.Fatalf("now with monotonic reading must compare as wall clock at an equal wall-only deadline") - } -} + 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) -func TestLateTimerFiring(t *testing.T) { - tests := []struct { - name string - final bool - beforeDl time.Duration - wantWarns int - wantFinals int - }{ - {"warning on resume inside window", false, 3 * time.Minute, 1, 0}, - {"warning promoted to final inside final window", false, time.Minute, 0, 1}, - {"warning skipped past deadline", false, -time.Minute, 0, 0}, - {"final on resume before deadline", true, time.Minute, 0, 1}, - {"final skipped past deadline", true, -time.Minute, 0, 0}, - } - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - r := &fakeRecorder{} - w := New(r) - defer w.Close() - - // The deadline is an hour out so the real timers never fire - // during the test; the late callback is invoked directly with an - // injected clock that simulates a resume near the deadline. - d := time.Now().Add(time.Hour).Round(0) - w.nowFn = func() time.Time { return d.Add(-tt.beforeDl) } - if err := w.Update(d); err != nil { - t.Fatalf("Update: %v", err) - } - - if tt.final { - w.fireFinal(d) - } else { - w.fire(d) - } - - events := r.snapshot() - if got := countWhere(events, event.isWarning); got != tt.wantWarns { - t.Fatalf("expected %d warning publishes, got %d: %+v", tt.wantWarns, got, events) - } - if got := countWhere(events, event.isFinalWarning); got != tt.wantFinals { - t.Fatalf("expected %d final-warning publishes, got %d: %+v", tt.wantFinals, got, events) - } - }) - } -} - -func TestPromotedFinalWarningIsNotRepeated(t *testing.T) { - r := &fakeRecorder{} - w := New(r) - defer w.Close() - - d := time.Now().Add(time.Hour).Round(0) - now := d.Add(-time.Minute) - w.nowFn = func() time.Time { return now } - if err := w.Update(d); err != nil { + deadline := start.Add(time.Hour).Round(0) + if err := w.Update(deadline); err != nil { t.Fatalf("Update: %v", err) } + return w, clock, deadline +} - w.fire(d) - // The final timer was suspended too, so it fires even later than the - // warning timer, here still just before the deadline. - now = d.Add(-30 * time.Second) - w.fireFinal(d) +func TestResumeInsideWarningWindowWarns(t *testing.T) { + r := &fakeRecorder{} + _, clock, deadline := newResumeWatcher(t, r) - events := r.snapshot() - if got := countWhere(events, event.isFinalWarning); got != 1 { - t.Fatalf("expected exactly 1 final-warning publish, got %d: %+v", got, events) + 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 got := countWhere(events, event.isWarning); got != 0 { - t.Fatalf("expected no regular warning publish, got %d: %+v", got, events) + if n := countWhere(events, event.isFinalWarning); n != 0 { + t.Fatalf("final-warning must wait for its own window, got %d: %+v", n, events) } } -func TestPromotionRespectsDismiss(t *testing.T) { +func TestResumeInsideFinalWindowSendsFinalWarningOnly(t *testing.T) { r := &fakeRecorder{} - w := New(r) - defer w.Close() + _, clock, deadline := newResumeWatcher(t, r) - d := time.Now().Add(time.Hour).Round(0) - w.nowFn = func() time.Time { return d.Add(-time.Minute) } - if err := w.Update(d); err != nil { - t.Fatalf("Update: %v", err) + 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() - w.fire(d) + clock.set(deadline.Add(-time.Minute)) - events := r.snapshot() - if got := countWhere(events, func(e event) bool { return e.kind == publish }); got != 0 { - t.Fatalf("expected no publish after dismiss, got %d: %+v", got, events) + 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 TestPromotionSkippedWhenFinalAlreadyFired(t *testing.T) { +func TestCloseStopsTheEvaluationLoop(t *testing.T) { r := &fakeRecorder{} - w := New(r) - defer w.Close() + w := newWatcherWithLeads(WarningLead, FinalWarningLead, r) + start := time.Now() + clock := newFakeClock(start) + w.nowFn = clock.now - // Both timers fall in the past after a long suspend and are dispatched - // with a zero delay, so the final callback can run before the warning one. - d := time.Now().Add(time.Hour).Round(0) - w.nowFn = func() time.Time { return d.Add(-time.Minute) } - if err := w.Update(d); err != nil { + deadline := start.Add(time.Hour).Round(0) + if err := w.Update(deadline); err != nil { t.Fatalf("Update: %v", err) } + w.Close() - w.fireFinal(d) - w.fire(d) + clock.set(deadline.Add(-5 * time.Minute)) + settle() - events := r.snapshot() - if got := countWhere(events, event.isFinalWarning); got != 1 { - t.Fatalf("expected exactly 1 final-warning publish, got %d: %+v", got, events) - } - if got := countWhere(events, event.isWarning); got != 0 { - t.Fatalf("expected no regular warning publish, got %d: %+v", got, events) + 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()) } } @@ -687,8 +722,8 @@ func TestDeadlineOnlyRecordsDeadlineWithoutWarnings(t *testing.T) { w := NewDeadlineOnly(r) defer w.Close() - // With the default leads this deadline would otherwise fire both - // timers on the next tick. + // 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) @@ -703,8 +738,11 @@ func TestDeadlineOnlyRecordsDeadlineWithoutWarnings(t *testing.T) { 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) } - if w.timer != nil || w.finalTimer != nil { - t.Fatal("expected no timers armed in deadline-only mode") + w.mu.Lock() + polling := w.stop != nil + w.mu.Unlock() + if polling { + t.Fatal("expected no evaluation loop in deadline-only mode") } } @@ -725,3 +763,157 @@ func TestDeadlineOnlyStillRejectsOutOfRangeDeadlines(t *testing.T) { 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") + } + } +}