feat: one Poll Lane per Site, replacing the shared pace (#100)

Each Site runs its own poll goroutine, paced by Rest and Gap from the
registry (sites.go) instead of the shared cooldown/interval/stagger/batch
config. Rest is enforced by the due query's WHERE clause; the Lane sleeps
its effective gap between fetches (hour/eligible, floored at 1s).

Lane-local failures: two refusals stop that Site for 15m, a lost browser
stops only the round's remaining browser Lanes, and cover heals moved to
background goroutines so a slow CDN cannot consume a Lane's gap. The five
LATEST_CHAPTER_POLL_* pace env vars are gone; only the kill switch
remains.
This commit is contained in:
2026-08-16 12:51:01 +07:00
parent ddbd57070d
commit 90ce14aec6
22 changed files with 3645 additions and 1884 deletions
+39 -16
View File
@@ -16,7 +16,6 @@ import (
"strconv"
"strings"
"github.com/jackc/pgx/v5/pgtype"
_ "github.com/jackc/pgx/v5/stdlib"
)
@@ -39,7 +38,7 @@ type Bookmark struct {
// origin once the bytes exist, and "" until they do — never a third-party
// address and never an address that 404s (ADR-0007). A client may still
// send this field and it is discarded on the way in; see Upsert.
Cover string `json:"cover"`
Cover string `json:"cover"`
LastChapter string `json:"last_chapter"`
LastChapterNum float64 `json:"last_chapter_num"`
LastChapterURL string `json:"last_chapter_url"`
@@ -897,10 +896,13 @@ func (s *Store) Delete(readerID int64, key string) error {
return nil
}
// DueForLatestCheck returns series whose server-side latest-chapter check has
// aged past the appropriate cutoff, ordered by how many bookmarks reference
// them (descending) then least-recently-checked first, at most limit of them.
// Browser-backed sites use browserCutoffMs; every other site uses cutoffMs.
// DueForLatestCheck returns one Site's series whose server-side
// latest-chapter check has aged past cutoffMs, ordered by how many bookmarks
// reference them (descending) then least-recently-checked first. One Site per
// query, because each Poll Lane asks for its own list: the query carries one
// Site and one cut-off instead of parallel lists (issue #100). There is no
// limit — the Lane's own gap paces the fetches, and the batch size that used
// to cap this query is gone with the shared pace.
//
// The reader_count ordering is the point of the split (ADR-0003): a series
// shared by several readers is fetched once per due cycle, and the popular
@@ -916,20 +918,17 @@ func (s *Store) Delete(readerID int64, key string) error {
// burns requests. Archived bookmarks still count — knowing what a shelved
// series is up to is the whole reason for archiving instead of deleting.
// A series with no bookmarks at all never appears: the join excludes it.
func (s *Store) DueForLatestCheck(cutoffMs, browserCutoffMs int64, browserSites []string, limit int) ([]Series, error) {
func (s *Store) DueForLatestCheck(site string, cutoffMs int64) ([]Series, error) {
rows, err := s.db.Query(`SELECT `+seriesColumns+`, COUNT(*) AS reader_count
FROM series s
JOIN bookmarks b ON b.site = s.site AND b.series_id = s.series_id
WHERE s.series_url <> ''
AND s.latest_checked_at <= CASE
WHEN s.site = ANY($3::text[]) THEN $2::bigint
ELSE $1::bigint
END
WHERE s.site = $1
AND s.series_url <> ''
AND s.latest_checked_at <= $2::bigint
GROUP BY s.site, s.series_id, s.title, s.series_url, s.cover,
s.kind, s.latest_chapter, s.latest_chapter_num, s.latest_checked_at
HAVING COUNT(*) FILTER (WHERE b.status <> 'finished') > 0
ORDER BY reader_count DESC, s.latest_checked_at ASC
LIMIT $4`, cutoffMs, browserCutoffMs, pgtype.FlatArray[string](browserSites), limit)
ORDER BY reader_count DESC, s.latest_checked_at ASC`, site, cutoffMs)
if err != nil {
return nil, fmt.Errorf("query due series: %w", err)
}
@@ -946,6 +945,30 @@ func (s *Store) DueForLatestCheck(cutoffMs, browserCutoffMs int64, browserSites
return out, rows.Err()
}
// EligibleSeriesCount returns how many of a Site's Series still have at least
// one bookmark outside the finished bucket. It is the denominator of the
// Lane's pace (issue #100): the effective gap is the smaller of the registry
// gap and one hour divided by this count, so Series that will never be Polled
// do not make the Lane faster than it needs to be, and counting every eligible
// Series rather than only those currently due keeps the pace steady — the
// single worst moment to be fastest is startup, when everything is due at
// once.
func (s *Store) EligibleSeriesCount(site string) (int, error) {
var n int
err := s.db.QueryRow(`SELECT COUNT(*) FROM (
SELECT 1
FROM series s
JOIN bookmarks b ON b.site = s.site AND b.series_id = s.series_id
WHERE s.site = $1
GROUP BY s.site, s.series_id
HAVING COUNT(*) FILTER (WHERE b.status <> 'finished') > 0
) e`, site).Scan(&n)
if err != nil {
return 0, fmt.Errorf("count eligible series %s: %w", site, err)
}
return n, nil
}
// MarkLatestChecked records that the server looked at a series at ts, whatever
// the look turned up. Marking a missing series is not an error: the row may
// have been orphaned while a fetch was in flight.
@@ -954,7 +977,7 @@ func (s *Store) DueForLatestCheck(cutoffMs, browserCutoffMs int64, browserSites
// out of the client-visible read path on purpose. PUT /bookmarks/{key} decodes
// a whole Bookmark from the client and Upsert writes every series column it
// knows about, so a userscript PUT — which has no idea this field exists —
// would write a zero and reset the cooldown, making the poller re-fetch that
// would write a zero and reset the rest, making the poller re-fetch that
// series every tick for as long as the user kept reading it.
func (s *Store) MarkLatestChecked(site, seriesID string, ts int64) error {
if _, err := s.db.Exec(
@@ -966,7 +989,7 @@ func (s *Store) MarkLatestChecked(site, seriesID string, ts int64) error {
}
// LatestCheckedAt reads the column MarkLatestChecked writes. It exists for
// tests outside this package (the poller's own tests assert on cooldown
// tests outside this package (the poller's own tests assert on rest
// bookkeeping) — see MarkLatestChecked for why the field stays off the
// client-visible row.
func (s *Store) LatestCheckedAt(site, seriesID string) (int64, error) {
+47 -12
View File
@@ -298,7 +298,7 @@ func TestDueForLatestCheck(t *testing.T) {
s := newTestStore(t)
seedForCheck(t, s, "asura:x", tt.seriesURL, tt.checkedAt)
due, err := s.DueForLatestCheck(now-hour, now-hour, nil, 10)
due, err := s.DueForLatestCheck("asura", now-hour)
if err != nil {
t.Fatalf("DueForLatestCheck: %v", err)
}
@@ -309,22 +309,25 @@ func TestDueForLatestCheck(t *testing.T) {
}
}
func TestDueForLatestCheckOldestFirstAndLimited(t *testing.T) {
func TestDueForLatestCheckOldestFirstAndScopedToSite(t *testing.T) {
s := newTestStore(t)
// Insert newest-checked first so a correct ORDER BY has to reverse it.
seedForCheck(t, s, "asura:c", "https://asurascans.com/comics/c", 300)
seedForCheck(t, s, "asura:b", "https://asurascans.com/comics/b", 200)
seedForCheck(t, s, "asura:a", "https://asurascans.com/comics/a", 100)
// A second Site's due series must not appear in asura's list: each Lane
// asks for one Site, and no Lane may see another's queue.
seedForCheck(t, s, "demonic:z", "https://demonicscans.org/manga/z", 0)
due, err := s.DueForLatestCheck(1000, 1000, nil, 2)
due, err := s.DueForLatestCheck("asura", 1000)
if err != nil {
t.Fatalf("DueForLatestCheck: %v", err)
}
if len(due) != 2 {
t.Fatalf("got %d rows, want 2 (limit)", len(due))
if len(due) != 3 {
t.Fatalf("got %d rows, want 3 (all of asura's, none of demonic's)", len(due))
}
if due[0].Key() != "asura:a" || due[1].Key() != "asura:b" {
t.Fatalf("got %q,%q; want asura:a,asura:b (oldest first)", due[0].Key(), due[1].Key())
if due[0].Key() != "asura:a" || due[1].Key() != "asura:b" || due[2].Key() != "asura:c" {
t.Fatalf("got %q,%q,%q; want asura:a,asura:b,asura:c (oldest first)", due[0].Key(), due[1].Key(), due[2].Key())
}
}
@@ -496,7 +499,7 @@ func TestDueForLatestCheckSkipsFinishedKeepsArchived(t *testing.T) {
}
}
due, err := store.DueForLatestCheck(time.Now().UnixMilli(), time.Now().UnixMilli(), nil, 10)
due, err := store.DueForLatestCheck("asura", time.Now().UnixMilli())
if err != nil {
t.Fatalf("DueForLatestCheck: %v", err)
}
@@ -512,6 +515,38 @@ func TestDueForLatestCheckSkipsFinishedKeepsArchived(t *testing.T) {
}
}
// The gap's denominator counts every Series the Lane will ever Poll: a
// finished Series must not make the Lane faster than it needs to be, and
// another Site's Series must not leak into this Site's count.
func TestEligibleSeriesCount(t *testing.T) {
store := newTestStore(t)
seedForCheck(t, store, "asura:reading", "https://asurascans.com/comics/reading", 0)
seedForCheck(t, store, "asura:archived", "https://asurascans.com/comics/archived", 0)
if _, err := store.Upsert(store.OwnerID(), Bookmark{
Key: "asura:finished", Site: "asura", SeriesID: "finished",
SeriesURL: "https://asurascans.com/comics/finished",
Status: StatusFinished, UpdatedAt: 1000,
}); err != nil {
t.Fatalf("seed finished: %v", err)
}
seedForCheck(t, store, "demonic:z", "https://demonicscans.org/manga/z", 0)
n, err := store.EligibleSeriesCount("asura")
if err != nil {
t.Fatalf("EligibleSeriesCount: %v", err)
}
if n != 2 {
t.Fatalf("eligible = %d, want 2 (finished excluded, demonic excluded)", n)
}
n, err = store.EligibleSeriesCount("demonic")
if err != nil {
t.Fatalf("EligibleSeriesCount(demonic): %v", err)
}
if n != 1 {
t.Fatalf("eligible(demonic) = %d, want 1", n)
}
}
func TestDisplayChapter(t *testing.T) {
cases := []struct {
name string
@@ -970,7 +1005,7 @@ func TestDueForLatestCheckOrdersByReaderCountThenAge(t *testing.T) {
seedSecondReader(t, s, "asura:pop:2", "asura", "pop", 1001)
seedForCheck(t, s, "asura:solo", "https://asurascans.com/comics/solo", 100)
due, err := s.DueForLatestCheck(1000, 1000, nil, 10)
due, err := s.DueForLatestCheck("asura", 1000)
if err != nil {
t.Fatalf("DueForLatestCheck: %v", err)
}
@@ -996,7 +1031,7 @@ func TestDueForLatestCheckExcludesOrphanSeries(t *testing.T) {
t.Fatalf("seed orphan series: %v", err)
}
due, err := s.DueForLatestCheck(1000, 1000, nil, 10)
due, err := s.DueForLatestCheck("asura", 1000)
if err != nil {
t.Fatalf("DueForLatestCheck: %v", err)
}
@@ -1334,7 +1369,7 @@ func TestTwoReadersShareOneSeriesWithIndependentProgress(t *testing.T) {
t.Fatalf("series rows = %d, want 1 shared row for two bookmarks", series)
}
due, err := s.DueForLatestCheck(time.Now().UnixMilli(), time.Now().UnixMilli(), nil, 10)
due, err := s.DueForLatestCheck("asura", time.Now().UnixMilli())
if err != nil {
t.Fatalf("DueForLatestCheck: %v", err)
}
@@ -1350,7 +1385,7 @@ func TestTwoReadersShareOneSeriesWithIndependentProgress(t *testing.T) {
if b, ok, err := s.Get(s.OwnerID(), "asura:solo"); err != nil || !ok || b.LastChapterNum != 200 {
t.Fatalf("owner's bookmark after the other's delete = %+v ok=%v err=%v, want it intact", b, ok, err)
}
due, err = s.DueForLatestCheck(time.Now().UnixMilli(), time.Now().UnixMilli(), nil, 10)
due, err = s.DueForLatestCheck("asura", time.Now().UnixMilli())
if err != nil {
t.Fatalf("DueForLatestCheck after delete: %v", err)
}