From 3211b5b42fea10f1705b16adb14dd7ba29f5080d Mon Sep 17 00:00:00 2001 From: Sulthan Zaki Date: Sun, 26 Jul 2026 15:40:17 +0700 Subject: [PATCH] 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. --- backend/latest.go | 162 ++++++++++++++++++++++++ backend/latest_test.go | 278 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 440 insertions(+) create mode 100644 backend/latest.go create mode 100644 backend/latest_test.go diff --git a/backend/latest.go b/backend/latest.go new file mode 100644 index 0000000..ddac2b4 --- /dev/null +++ b/backend/latest.go @@ -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) +} diff --git a/backend/latest_test.go b/backend/latest_test.go new file mode 100644 index 0000000..19356c1 --- /dev/null +++ b/backend/latest_test.go @@ -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) + } +}