feat: background poller for latest published chapter
Ticker goroutine reads bookmarks past their per-bookmark cooldown, fetches the series page, and writes latest_chapter through Get+Upsert so updated_at never moves and the list never reorders. The row is stamped before the fetch so a broken series waits out a cooldown instead of retrying every tick.
This commit is contained in:
@@ -0,0 +1,162 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"log"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// fetcher retrieves a series page. It exists as an interface so tests can inject
|
||||||
|
// a fake: nothing in the test suite may touch the network or the TLS client.
|
||||||
|
type fetcher interface {
|
||||||
|
Get(ctx context.Context, url string) (body string, status int, err error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// latestPoller re-checks each bookmarked series' newest published chapter on a
|
||||||
|
// schedule, independent of the userscript's own in-browser checks. The two run
|
||||||
|
// in parallel and report the same observable fact, so whichever writes last wins
|
||||||
|
// and neither needs to know about the other.
|
||||||
|
//
|
||||||
|
// Two clocks, deliberately independent:
|
||||||
|
//
|
||||||
|
// - interval is how often this goroutine wakes up and looks.
|
||||||
|
// - cooldown is how long one bookmark rests since its own last check.
|
||||||
|
//
|
||||||
|
// Only the cooldown is per bookmark, and it is enforced by the WHERE clause in
|
||||||
|
// DueForLatestCheck rather than by any timer. Shortening interval therefore
|
||||||
|
// cannot shorten anyone's cooldown; it only makes the poller wake up and find
|
||||||
|
// nothing due more often.
|
||||||
|
type latestPoller struct {
|
||||||
|
store *Store
|
||||||
|
fetch fetcher
|
||||||
|
now func() time.Time // injected so tests can freeze it
|
||||||
|
cooldown time.Duration
|
||||||
|
interval time.Duration
|
||||||
|
stagger time.Duration
|
||||||
|
batch int
|
||||||
|
}
|
||||||
|
|
||||||
|
// Run polls until ctx is cancelled.
|
||||||
|
//
|
||||||
|
// runOnce is called synchronously, so a batch that overruns the tick delays the
|
||||||
|
// next one instead of stacking a second batch on top of it. That is the intended
|
||||||
|
// failure mode for a misconfigured batch x stagger: a slower cadence, never
|
||||||
|
// concurrent fetch storms.
|
||||||
|
func (p *latestPoller) Run(ctx context.Context) {
|
||||||
|
log.Printf("latest-chapter poller: interval=%s cooldown=%s batch=%d stagger=%s",
|
||||||
|
p.interval, p.cooldown, p.batch, p.stagger)
|
||||||
|
t := time.NewTicker(p.interval)
|
||||||
|
defer t.Stop()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
log.Println("latest-chapter poller: stopped")
|
||||||
|
return
|
||||||
|
case <-t.C:
|
||||||
|
p.runOnce(ctx)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// runOnce processes one batch of due bookmarks.
|
||||||
|
func (p *latestPoller) runOnce(ctx context.Context) {
|
||||||
|
cutoff := p.now().Add(-p.cooldown).UnixMilli()
|
||||||
|
due, err := p.store.DueForLatestCheck(cutoff, p.batch)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("latest poll: due query: %v", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if len(due) == 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
checked := 0
|
||||||
|
for i, b := range due {
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
// Staggered rather than fired together: a burst of simultaneous requests
|
||||||
|
// from one server IP is the traffic shape most likely to move that IP's
|
||||||
|
// bot score. This is the server-side analogue of the userscript's "one
|
||||||
|
// series per navigation ... indistinguishable from browsing" (L455-456).
|
||||||
|
if i > 0 && p.stagger > 0 {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
case <-time.After(p.stagger):
|
||||||
|
}
|
||||||
|
}
|
||||||
|
p.checkOne(ctx, b)
|
||||||
|
checked++
|
||||||
|
}
|
||||||
|
// due vs checked is how you tell which constraint is binding: ticks that
|
||||||
|
// report due=0 mean the cooldown is the limit, ticks that report due==batch
|
||||||
|
// every time mean throughput is.
|
||||||
|
log.Printf("latest poll: due=%d checked=%d", len(due), checked)
|
||||||
|
}
|
||||||
|
|
||||||
|
// checkOne re-checks one series. Every failure path here is "log and move on":
|
||||||
|
// the poller is a best-effort enhancement, and no single bad series may stall a
|
||||||
|
// batch or take down the process.
|
||||||
|
func (p *latestPoller) checkOne(ctx context.Context, b Bookmark) {
|
||||||
|
defer func() {
|
||||||
|
if r := recover(); r != nil {
|
||||||
|
log.Printf("latest poll %q: recovered from panic: %v", b.Key, r)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
// Stamped before the fetch, not after, so an error, a timeout, or a shutdown
|
||||||
|
// mid-request still consumes the cooldown. Otherwise a renamed or deleted
|
||||||
|
// series would be retried on every single tick forever. The userscript
|
||||||
|
// stamps in the same order and for the same reason (L471-473).
|
||||||
|
if err := p.store.MarkLatestChecked(b.Key, p.now().UnixMilli()); err != nil {
|
||||||
|
log.Printf("latest poll %q: mark checked: %v", b.Key, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
body, status, err := p.fetch.Get(ctx, b.SeriesURL)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("latest poll %q: fetch %s: %v", b.Key, b.SeriesURL, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if status != 200 {
|
||||||
|
log.Printf("latest poll %q: fetch %s: status %d", b.Key, b.SeriesURL, status)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
latest, ok := latestChapterFrom(b.Site, b.SeriesURL, body)
|
||||||
|
if !ok {
|
||||||
|
// Most likely a challenge page or a layout change. Either way the row is
|
||||||
|
// already stamped, so this waits out a cooldown instead of hot-looping.
|
||||||
|
log.Printf("latest poll %q: no chapter links in %d bytes", b.Key, len(body))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// Re-read: the row may have been updated or deleted while the fetch was in
|
||||||
|
// flight, and writing b back wholesale would undo that.
|
||||||
|
cur, found, err := p.store.Get(b.Key)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("latest poll %q: reread: %v", b.Key, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if !found {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
// Equality, not >, mirroring the userscript (L427): a site that retracts a
|
||||||
|
// chapter should correct the stored number downward.
|
||||||
|
if cur.LatestChapterNum != nil && *cur.LatestChapterNum == latest.Num {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
num := latest.Num
|
||||||
|
cur.LatestChapter = latest.Label
|
||||||
|
cur.LatestChapterNum = &num
|
||||||
|
// A candidate only. last_chapter_num is untouched, so the CASE in Upsert
|
||||||
|
// keeps the stored updated_at and the bookmark list does not reorder.
|
||||||
|
cur.UpdatedAt = p.now().UnixMilli()
|
||||||
|
if _, err := p.store.Upsert(cur); err != nil {
|
||||||
|
log.Printf("latest poll %q: upsert: %v", b.Key, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
log.Printf("latest poll %q: latest is now %s", b.Key, latest.Label)
|
||||||
|
}
|
||||||
@@ -0,0 +1,278 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"sync"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// fakeFetcher stands in for the network. Every poller test uses it, so nothing
|
||||||
|
// in this file can reach tls-client or a real site.
|
||||||
|
type fakeFetcher struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
calls []string
|
||||||
|
body string
|
||||||
|
status int
|
||||||
|
err error
|
||||||
|
// perURL overrides body/status/err for specific URLs.
|
||||||
|
perURL map[string]fakeResponse
|
||||||
|
}
|
||||||
|
|
||||||
|
type fakeResponse struct {
|
||||||
|
body string
|
||||||
|
status int
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeFetcher) Get(ctx context.Context, url string) (string, int, error) {
|
||||||
|
f.mu.Lock()
|
||||||
|
f.calls = append(f.calls, url)
|
||||||
|
f.mu.Unlock()
|
||||||
|
if r, ok := f.perURL[url]; ok {
|
||||||
|
return r.body, r.status, r.err
|
||||||
|
}
|
||||||
|
return f.body, f.status, f.err
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeFetcher) callCount() int {
|
||||||
|
f.mu.Lock()
|
||||||
|
defer f.mu.Unlock()
|
||||||
|
return len(f.calls)
|
||||||
|
}
|
||||||
|
|
||||||
|
// newTestPoller wires a poller with a frozen clock and no stagger, so tests run
|
||||||
|
// instantly and deterministically.
|
||||||
|
func newTestPoller(t *testing.T, s *Store, f fetcher, at time.Time) *latestPoller {
|
||||||
|
t.Helper()
|
||||||
|
return &latestPoller{
|
||||||
|
store: s,
|
||||||
|
fetch: f,
|
||||||
|
now: func() time.Time { return at },
|
||||||
|
cooldown: time.Hour,
|
||||||
|
interval: 10 * time.Minute,
|
||||||
|
stagger: 0,
|
||||||
|
batch: 14,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRunOnceRecordsLatestChapter(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
const url = "https://asurascans.com/comics/chronicles-of-the-demon-faction-f886a8af"
|
||||||
|
seedForCheck(t, s, "asura:chronicles-of-the-demon-faction-f886a8af", url, 0)
|
||||||
|
|
||||||
|
now := time.UnixMilli(5_000_000)
|
||||||
|
f := &fakeFetcher{body: asuraSeriesFixture, status: 200}
|
||||||
|
newTestPoller(t, s, f, now).runOnce(context.Background())
|
||||||
|
|
||||||
|
b, ok, err := s.Get("asura:chronicles-of-the-demon-faction-f886a8af")
|
||||||
|
if err != nil || !ok {
|
||||||
|
t.Fatalf("Get: %v ok=%v", err, ok)
|
||||||
|
}
|
||||||
|
if b.LatestChapterNum == nil || *b.LatestChapterNum != 181 {
|
||||||
|
t.Fatalf("LatestChapterNum = %v, want 181", b.LatestChapterNum)
|
||||||
|
}
|
||||||
|
if b.LatestChapter != "Chapter 181" {
|
||||||
|
t.Fatalf("LatestChapter = %q, want %q", b.LatestChapter, "Chapter 181")
|
||||||
|
}
|
||||||
|
if got := readLatestCheckedAt(t, s, "asura:chronicles-of-the-demon-faction-f886a8af"); got != now.UnixMilli() {
|
||||||
|
t.Fatalf("latest_checked_at = %d, want %d", got, now.UnixMilli())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The whole point of the updated_at CASE in Upsert: a newly published chapter is
|
||||||
|
// not reading progress and must not move the series up the list.
|
||||||
|
func TestRunOnceDoesNotReorderList(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
const url = "https://asurascans.com/comics/chronicles-of-the-demon-faction-f886a8af"
|
||||||
|
const key = "asura:chronicles-of-the-demon-faction-f886a8af"
|
||||||
|
|
||||||
|
// "other" is the most recently read, so it must stay at the top of List().
|
||||||
|
if _, err := s.Upsert(Bookmark{
|
||||||
|
Key: "asura:other", Site: "asura", SeriesID: "other",
|
||||||
|
SeriesURL: "https://asurascans.com/comics/other", UpdatedAt: 9_000_000,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("seed other: %v", err)
|
||||||
|
}
|
||||||
|
seedForCheck(t, s, key, url, 0)
|
||||||
|
before, _, err := s.Get(key)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Get before: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
f := &fakeFetcher{body: asuraSeriesFixture, status: 200}
|
||||||
|
newTestPoller(t, s, f, time.UnixMilli(9_999_999)).runOnce(context.Background())
|
||||||
|
|
||||||
|
after, _, err := s.Get(key)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Get after: %v", err)
|
||||||
|
}
|
||||||
|
if after.UpdatedAt != before.UpdatedAt {
|
||||||
|
t.Fatalf("updated_at moved from %d to %d on a latest-chapter bump",
|
||||||
|
before.UpdatedAt, after.UpdatedAt)
|
||||||
|
}
|
||||||
|
list, err := s.List()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("List: %v", err)
|
||||||
|
}
|
||||||
|
if list[0].Key != "asura:other" {
|
||||||
|
t.Fatalf("list reordered: head is %q, want asura:other", list[0].Key)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A failed fetch must still consume the cooldown, or a renamed series gets
|
||||||
|
// retried on every tick forever.
|
||||||
|
func TestRunOnceMarksCheckedOnFailure(t *testing.T) {
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
resp fakeResponse
|
||||||
|
}{
|
||||||
|
{"network error", fakeResponse{err: errors.New("dial tcp: refused")}},
|
||||||
|
{"non-200", fakeResponse{body: "nope", status: 503}},
|
||||||
|
{"challenge page", fakeResponse{body: challengeFixture, status: 200}},
|
||||||
|
{"empty body", fakeResponse{body: "", status: 200}},
|
||||||
|
}
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
const url = "https://asurascans.com/comics/x"
|
||||||
|
seedForCheck(t, s, "asura:x", url, 0)
|
||||||
|
|
||||||
|
now := time.UnixMilli(7_000_000)
|
||||||
|
f := &fakeFetcher{perURL: map[string]fakeResponse{url: tt.resp}}
|
||||||
|
newTestPoller(t, s, f, now).runOnce(context.Background())
|
||||||
|
|
||||||
|
if got := readLatestCheckedAt(t, s, "asura:x"); got != now.UnixMilli() {
|
||||||
|
t.Fatalf("latest_checked_at = %d, want %d", got, now.UnixMilli())
|
||||||
|
}
|
||||||
|
b, _, err := s.Get("asura:x")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Get: %v", err)
|
||||||
|
}
|
||||||
|
if b.LatestChapterNum != nil {
|
||||||
|
t.Fatalf("LatestChapterNum = %v, want nil on a failed check", *b.LatestChapterNum)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRunOnceRespectsBatchLimit(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
for i := 0; i < 20; i++ {
|
||||||
|
key := "asura:s" + string(rune('a'+i))
|
||||||
|
seedForCheck(t, s, key, "https://asurascans.com/comics/"+key, 0)
|
||||||
|
}
|
||||||
|
|
||||||
|
f := &fakeFetcher{body: "", status: 200}
|
||||||
|
p := newTestPoller(t, s, f, time.UnixMilli(5_000_000))
|
||||||
|
p.batch = 5
|
||||||
|
p.runOnce(context.Background())
|
||||||
|
|
||||||
|
if got := f.callCount(); got != 5 {
|
||||||
|
t.Fatalf("fetched %d series, want 5 (batch limit)", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// One unreachable series must not abandon the rest of the batch.
|
||||||
|
func TestRunOnceOneBadSeriesDoesNotStallBatch(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
keys := []string{"asura:a", "asura:b", "asura:c", "asura:d", "asura:e"}
|
||||||
|
for _, k := range keys {
|
||||||
|
seedForCheck(t, s, k, "https://asurascans.com/comics/"+k, 0)
|
||||||
|
}
|
||||||
|
|
||||||
|
now := time.UnixMilli(6_000_000)
|
||||||
|
f := &fakeFetcher{
|
||||||
|
body: "", status: 200,
|
||||||
|
perURL: map[string]fakeResponse{
|
||||||
|
"https://asurascans.com/comics/asura:b": {err: errors.New("boom")},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
newTestPoller(t, s, f, now).runOnce(context.Background())
|
||||||
|
|
||||||
|
if got := f.callCount(); got != 5 {
|
||||||
|
t.Fatalf("fetched %d series, want all 5 attempted", got)
|
||||||
|
}
|
||||||
|
for _, k := range keys {
|
||||||
|
if got := readLatestCheckedAt(t, s, k); got != now.UnixMilli() {
|
||||||
|
t.Fatalf("%s latest_checked_at = %d, want %d", k, got, now.UnixMilli())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The cooldown is enforced by the due query, so a second immediate pass must do
|
||||||
|
// nothing at all — this is what makes the tick interval independent of it.
|
||||||
|
func TestRunOnceHonoursCooldownAcrossPasses(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
const url = "https://asurascans.com/comics/x"
|
||||||
|
seedForCheck(t, s, "asura:x", url, 0)
|
||||||
|
|
||||||
|
now := time.UnixMilli(8_000_000)
|
||||||
|
f := &fakeFetcher{body: asuraSeriesFixture, status: 200}
|
||||||
|
p := newTestPoller(t, s, f, now)
|
||||||
|
|
||||||
|
p.runOnce(context.Background())
|
||||||
|
if got := f.callCount(); got != 1 {
|
||||||
|
t.Fatalf("first pass fetched %d, want 1", got)
|
||||||
|
}
|
||||||
|
// Same instant, and again 59 minutes later: both inside the 1h cooldown.
|
||||||
|
p.runOnce(context.Background())
|
||||||
|
p.now = func() time.Time { return now.Add(59 * time.Minute) }
|
||||||
|
p.runOnce(context.Background())
|
||||||
|
if got := f.callCount(); got != 1 {
|
||||||
|
t.Fatalf("fetched %d times inside the cooldown, want 1", got)
|
||||||
|
}
|
||||||
|
// Past the cooldown, it is due again.
|
||||||
|
p.now = func() time.Time { return now.Add(61 * time.Minute) }
|
||||||
|
p.runOnce(context.Background())
|
||||||
|
if got := f.callCount(); got != 2 {
|
||||||
|
t.Fatalf("fetched %d times after the cooldown, want 2", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A site that retracts a chapter should correct the stored number downward,
|
||||||
|
// mirroring the userscript's equality check (L427) rather than a >.
|
||||||
|
func TestRunOnceCorrectsDownward(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
const url = "https://demonicscans.org/manga/Catastrophic-Necromancer"
|
||||||
|
const key = "demonic:Catastrophic-Necromancer"
|
||||||
|
|
||||||
|
high := 400.0
|
||||||
|
if _, err := s.Upsert(Bookmark{
|
||||||
|
Key: key, Site: "demonic", SeriesID: "Catastrophic-Necromancer",
|
||||||
|
SeriesURL: url, LatestChapter: "Chapter 400", LatestChapterNum: &high,
|
||||||
|
UpdatedAt: 1000,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("seed: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
f := &fakeFetcher{body: demonicSeriesFixture, status: 200}
|
||||||
|
newTestPoller(t, s, f, time.UnixMilli(5_000_000)).runOnce(context.Background())
|
||||||
|
|
||||||
|
b, _, err := s.Get(key)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Get: %v", err)
|
||||||
|
}
|
||||||
|
if b.LatestChapterNum == nil || *b.LatestChapterNum != 296 {
|
||||||
|
t.Fatalf("LatestChapterNum = %v, want 296", b.LatestChapterNum)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A cancelled context must abandon the batch rather than run it to completion.
|
||||||
|
func TestRunOnceStopsOnCancelledContext(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
for _, k := range []string{"asura:a", "asura:b", "asura:c"} {
|
||||||
|
seedForCheck(t, s, k, "https://asurascans.com/comics/"+k, 0)
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
cancel()
|
||||||
|
|
||||||
|
f := &fakeFetcher{body: "", status: 200}
|
||||||
|
newTestPoller(t, s, f, time.UnixMilli(5_000_000)).runOnce(ctx)
|
||||||
|
|
||||||
|
if got := f.callCount(); got != 0 {
|
||||||
|
t.Fatalf("fetched %d series with a cancelled context, want 0", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user