This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git


The following commit(s) were added to refs/heads/main by this push:
     new 1507bc010a test(auth): extend UserActivityTracker unit test coverage 
(#6879)
1507bc010a is described below

commit 1507bc010af3aa0a56e10d065e7320a3e273b9d0
Author: Xinyuan Lin <[email protected]>
AuthorDate: Fri Jul 24 22:17:18 2026 -0700

    test(auth): extend UserActivityTracker unit test coverage (#6879)
    
    ### What changes were proposed in this PR?
    
    Extends `UserActivityTrackerSpec` (22 tests total) for
    `UserActivityTracker` (~29% covered). The class ctor is fully
    injectable, so the tests use a fake clock, a direct executor and a
    recording upsert function to stay deterministic:
    - `markActive`: null-uid guard, first-call upsert, suppression while
    inside the write interval, re-upsert after it elapses, and an upsert
    failure being swallowed rather than propagated;
    - `evictStale`: entries older than 2x the write interval evicted while
    fresher ones are retained;
    - `cooldownSize` bookkeeping across those transitions.
    
    No source changes.
    
    ### Any related issues, documentation, discussions?
    
    Closes #6876.
    
    ### How was this PR tested?
    
    `sbt -java-home <jbr-17> "Auth/testOnly *UserActivityTrackerSpec"` -> 22
    succeeded, 0 failed. `Test/scalafmtCheck` + `Test/scalafix --check`
    clean.
    
    ### Was this PR authored or co-authored using generative AI tooling?
    
    Generated-by: Claude Code (Opus 4.8 [1M context])
---
 .../texera/auth/UserActivityTrackerSpec.scala      | 210 ++++++++++++++++++++-
 1 file changed, 209 insertions(+), 1 deletion(-)

diff --git 
a/common/auth/src/test/scala/org/apache/texera/auth/UserActivityTrackerSpec.scala
 
b/common/auth/src/test/scala/org/apache/texera/auth/UserActivityTrackerSpec.scala
index 5ffd402d9a..08fcad2097 100644
--- 
a/common/auth/src/test/scala/org/apache/texera/auth/UserActivityTrackerSpec.scala
+++ 
b/common/auth/src/test/scala/org/apache/texera/auth/UserActivityTrackerSpec.scala
@@ -23,7 +23,12 @@ import org.scalatest.flatspec.AnyFlatSpec
 import org.scalatest.matchers.should.Matchers
 
 import java.time.{Duration, Instant}
-import java.util.concurrent.{ConcurrentLinkedQueue, Executor}
+import java.util.concurrent.{
+  ConcurrentLinkedQueue,
+  CountDownLatch,
+  Executor,
+  RejectedExecutionException
+}
 import java.util.concurrent.atomic.AtomicReference
 
 class UserActivityTrackerSpec extends AnyFlatSpec with Matchers {
@@ -168,4 +173,207 @@ class UserActivityTrackerSpec extends AnyFlatSpec with 
Matchers {
 
     noException should be thrownBy tracker.evictStale()
   }
+
+  it should "not create a cooldown entry for a null uid" in {
+    val recorder = new Recorder
+    val clock = new 
AtomicReference[Instant](Instant.parse("2026-01-01T00:00:00Z"))
+    val tracker = makeTracker(Duration.ofMinutes(5), recorder, clock)
+
+    tracker.markActive(null)
+
+    tracker.cooldownSize shouldBe 0
+  }
+
+  it should "keep exactly one cooldown entry per uid across repeated calls" in 
{
+    val recorder = new Recorder
+    val t0 = Instant.parse("2026-01-01T00:00:00Z")
+    val clock = new AtomicReference[Instant](t0)
+    val tracker = makeTracker(Duration.ofMinutes(5), recorder, clock)
+
+    tracker.markActive(42)
+    clock.set(t0.plus(Duration.ofMinutes(1)))
+    tracker.markActive(42)
+    clock.set(t0.plus(Duration.ofMinutes(6)))
+    tracker.markActive(42)
+
+    recorder.calls.size shouldBe 2
+    tracker.cooldownSize shouldBe 1
+  }
+
+  it should "stamp each upsert with the clock value at claim time" in {
+    val recorder = new Recorder
+    val t0 = Instant.parse("2026-01-01T00:00:00Z")
+    val t1 = t0.plus(Duration.ofMinutes(7))
+    val clock = new AtomicReference[Instant](t0)
+    val tracker = makeTracker(Duration.ofMinutes(5), recorder, clock)
+
+    tracker.markActive(42)
+    // a suppressed call in between must not change the recorded timestamps
+    clock.set(t0.plus(Duration.ofMinutes(3)))
+    tracker.markActive(42)
+    clock.set(t1)
+    tracker.markActive(42)
+
+    recorder.calls.poll() shouldBe ((42, t0))
+    recorder.calls.poll() shouldBe ((42, t1))
+    recorder.calls.size shouldBe 0
+  }
+
+  it should "evict only the stale entries and retain the fresh ones" in {
+    val recorder = new Recorder
+    val t0 = Instant.parse("2026-01-01T00:00:00Z")
+    val clock = new AtomicReference[Instant](t0)
+    val tracker = makeTracker(Duration.ofMinutes(5), recorder, clock)
+
+    tracker.markActive(1)
+    clock.set(t0.plus(Duration.ofMinutes(8)))
+    tracker.markActive(2)
+    tracker.cooldownSize shouldBe 2
+
+    // cutoff = t0 + 11min - 10min = t0 + 1min: uid 1 (t0) is stale, uid 2 
(t0+8min) is not
+    clock.set(t0.plus(Duration.ofMinutes(11)))
+    tracker.evictStale()
+    tracker.cooldownSize shouldBe 1
+
+    // uid 2 is still in cooldown, so it must not produce a second write
+    tracker.markActive(2)
+    recorder.calls.size shouldBe 2
+  }
+
+  it should "retain an entry that is exactly 2 * writeInterval old" in {
+    val recorder = new Recorder
+    val t0 = Instant.parse("2026-01-01T00:00:00Z")
+    val clock = new AtomicReference[Instant](t0)
+    val tracker = makeTracker(Duration.ofMinutes(5), recorder, clock)
+
+    tracker.markActive(1)
+
+    // cutoff == t0 exactly; the boundary entry is not "before" the cutoff
+    clock.set(t0.plus(Duration.ofMinutes(10)))
+    tracker.evictStale()
+    tracker.cooldownSize shouldBe 1
+
+    // one millisecond past the boundary it is evicted
+    clock.set(t0.plus(Duration.ofMinutes(10)).plusMillis(1))
+    tracker.evictStale()
+    tracker.cooldownSize shouldBe 0
+  }
+
+  it should "tolerate evictStale on an empty tracker" in {
+    val recorder = new Recorder
+    val clock = new 
AtomicReference[Instant](Instant.parse("2026-01-01T00:00:00Z"))
+    val tracker = makeTracker(Duration.ofMinutes(5), recorder, clock)
+
+    noException should be thrownBy tracker.evictStale()
+
+    tracker.cooldownSize shouldBe 0
+    recorder.calls.size shouldBe 0
+  }
+
+  it should "re-claim a uid whose entry was evicted" in {
+    val recorder = new Recorder
+    val t0 = Instant.parse("2026-01-01T00:00:00Z")
+    val clock = new AtomicReference[Instant](t0)
+    val tracker = makeTracker(Duration.ofMinutes(5), recorder, clock)
+
+    tracker.markActive(1)
+    clock.set(t0.plus(Duration.ofMinutes(11)))
+    tracker.evictStale()
+    tracker.cooldownSize shouldBe 0
+
+    tracker.markActive(1)
+
+    recorder.calls.size shouldBe 2
+    tracker.cooldownSize shouldBe 1
+  }
+
+  it should "perform exactly one upsert when many threads race on the same 
uid" in {
+    val recorder = new Recorder
+    val clock = new 
AtomicReference[Instant](Instant.parse("2026-01-01T00:00:00Z"))
+    val tracker = makeTracker(Duration.ofMinutes(5), recorder, clock)
+
+    val start = new CountDownLatch(1)
+    val threads = (1 to 16).map(_ =>
+      new Thread(() => {
+        start.await()
+        tracker.markActive(99)
+      })
+    )
+    threads.foreach(_.start())
+    start.countDown()
+    // bounded join: a stuck thread fails the test instead of hanging the suite
+    threads.foreach(_.join(5000))
+    threads.zipWithIndex.foreach {
+      case (t, i) =>
+        withClue(s"thread $i did not finish within 5s: ")(t.isAlive shouldBe 
false)
+    }
+
+    // the CAS claim lets a single caller through; the rest are dropped
+    recorder.calls.size shouldBe 1
+    tracker.cooldownSize shouldBe 1
+  }
+
+  it should "upsert on every call when the write interval is zero" in {
+    val recorder = new Recorder
+    val clock = new 
AtomicReference[Instant](Instant.parse("2026-01-01T00:00:00Z"))
+    val tracker = makeTracker(Duration.ZERO, recorder, clock)
+
+    tracker.markActive(42)
+    tracker.markActive(42)
+    tracker.markActive(42)
+
+    recorder.calls.size shouldBe 3
+  }
+
+  it should "suppress the upsert when the clock moves backwards" in {
+    val recorder = new Recorder
+    val t0 = Instant.parse("2026-01-01T00:00:00Z")
+    val clock = new AtomicReference[Instant](t0)
+    val tracker = makeTracker(Duration.ofMinutes(5), recorder, clock)
+
+    tracker.markActive(42)
+    clock.set(t0.minus(Duration.ofMinutes(30)))
+    tracker.markActive(42)
+
+    // a negative elapsed time compares below the interval, so nothing is 
written
+    recorder.calls.size shouldBe 1
+    tracker.cooldownSize shouldBe 1
+  }
+
+  it should "keep the cooldown claim when the executor drops the write" in {
+    val recorder = new Recorder
+    val t0 = Instant.parse("2026-01-01T00:00:00Z")
+    val clock = new AtomicReference[Instant](t0)
+    val dropping: Executor = (_: Runnable) => ()
+    val tracker =
+      new UserActivityTracker(Duration.ofMinutes(5), recorder.upsert, 
dropping, () => clock.get())
+
+    tracker.markActive(42)
+    tracker.cooldownSize shouldBe 1
+
+    // the claim stands even though the write never ran, so the cooldown still 
holds
+    clock.set(t0.plus(Duration.ofMinutes(2)))
+    tracker.markActive(42)
+
+    recorder.calls.size shouldBe 0
+    tracker.cooldownSize shouldBe 1
+  }
+
+  it should "swallow a rejection thrown by the executor" in {
+    val recorder = new Recorder
+    val clock = new 
AtomicReference[Instant](Instant.parse("2026-01-01T00:00:00Z"))
+    val rejecting: Executor = (_: Runnable) => throw new 
RejectedExecutionException("queue full")
+    val tracker =
+      new UserActivityTracker(Duration.ofMinutes(5), recorder.upsert, 
rejecting, () => clock.get())
+
+    noException should be thrownBy tracker.markActive(42)
+
+    recorder.calls.size shouldBe 0
+    // the slot was claimed before the dispatch attempt failed
+    tracker.cooldownSize shouldBe 1
+  }
+
+  "UserActivityTracker singleton" should "treat a null uid as a no-op without 
touching the DB" in {
+    noException should be thrownBy UserActivityTracker.markActive(null)
+  }
 }

Reply via email to