A newly connected account showed no videos until the next 2h scheduled pass — the gap that made onboarding look broken (a second user connected, saw nothing, read as failure). The connect callback now fires an out-of-band discovery pass for the connecting user, so videos appear promptly. Concurrency: scheduled and connect-triggered passes share one lock (serialize), preserving the single-fetcher invariant (ADR-018). A trigger interleaves between the scheduler's per-user passes rather than fetching concurrently or waiting for a whole pass. The trigger runs on the server ctx (survives the redirect) and is non-blocking for the request goroutine. Scope: connect-trigger only. The optional login-refresh / "Discover now" button from #6 are intentionally not built — an unconditional login hook risks 429 storms (per the ticket's own recommendation); defer until wanted. TDD: TestCallbackTriggersDiscovery, TestSerializeRunsOneAtATime, TestDiscoveryTriggerEnqueueRunsUser; new BDD scenario mapped. Refs #6 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
63 lines
1.6 KiB
Go
63 lines
1.6 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"sync"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/stretchr/testify/require"
|
|
|
|
"gitea.d-ma.be/mathias/tapir/internal/runner"
|
|
)
|
|
|
|
// serialize must guarantee at most one discovery pass runs at a time, so a
|
|
// connect-triggered pass never fetches concurrently with the scheduler.
|
|
func TestSerializeRunsOneAtATime(t *testing.T) {
|
|
var active, maxActive int32
|
|
run := func(_ context.Context, _ string) (runner.Stats, error) {
|
|
n := atomic.AddInt32(&active, 1)
|
|
for { // record the high-water mark of concurrent runs
|
|
m := atomic.LoadInt32(&maxActive)
|
|
if n <= m || atomic.CompareAndSwapInt32(&maxActive, m, n) {
|
|
break
|
|
}
|
|
}
|
|
time.Sleep(2 * time.Millisecond)
|
|
atomic.AddInt32(&active, -1)
|
|
return runner.Stats{}, nil
|
|
}
|
|
|
|
s := serialize(&sync.Mutex{}, run)
|
|
var wg sync.WaitGroup
|
|
for i := 0; i < 20; i++ {
|
|
wg.Add(1)
|
|
go func(i int) { defer wg.Done(); _, _ = s(context.Background(), fmt.Sprintf("u%d", i)) }(i)
|
|
}
|
|
wg.Wait()
|
|
|
|
require.Equal(t, int32(1), atomic.LoadInt32(&maxActive),
|
|
"serialize must run at most one pass at a time")
|
|
}
|
|
|
|
// Enqueue runs the user's pass out-of-band (non-blocking) on the trigger's ctx.
|
|
func TestDiscoveryTriggerEnqueueRunsUser(t *testing.T) {
|
|
done := make(chan string, 1)
|
|
run := func(_ context.Context, userID string) (runner.Stats, error) {
|
|
done <- userID
|
|
return runner.Stats{}, nil
|
|
}
|
|
tr := &discoveryTrigger{ctx: context.Background(), run: run, log: quietLog()}
|
|
|
|
tr.Enqueue("u1")
|
|
|
|
select {
|
|
case got := <-done:
|
|
require.Equal(t, "u1", got)
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("Enqueue did not run the user's pass")
|
|
}
|
|
}
|