Closes #100. Each Site runs its own Poll Lane: an independent goroutine with its own rest and pace from the registry (`backend/internal/latest/sites.go`), replacing the shared cooldown/interval/stagger/batch configuration. Rest (1h, all six Sites including the browser trio) is enforced by the due query's WHERE clause; the Lane sleeps its effective gap between fetches — the registry 10s, or rest/eligible when a Site holds enough Series, floored at 1s with a Site-naming warning when the floor engages. Lane-local failure handling: - Two challenge-held results stop that Site's Lane for 15m; the probes keep their stamp, untried Series stay due. - A lost browser sets a shared Poller flag: the other browser Lanes skip their passes for the same 15m (no stamp-per-pass-per-Lane on a dead tab), then decay and probe again. - Browser wake gate preserved (5 due, or one waiting 15m, ADR-0005); one tab shared by the three browser Sites; "browser lane behind by X" logged every pass. - Cover work (healing a stored source URL and filling a blank from the series page) runs in the background so a slow CDN cannot consume a Lane's gap. Removed: `LATEST_CHAPTER_POLL_{COOLDOWN,BROWSER_COOLDOWN,INTERVAL,BATCH,STAGGER}` and the 6h browser rest. Only `LATEST_CHAPTER_POLL_ENABLED` remains; DEPLOY.md documents the exact `.env` edit. ADR-0010 records the decisions. Reviewed-on: #106 Co-authored-by: Sulthan Zaki <sultankiki05@gmail.com> Co-committed-by: Sulthan Zaki <sultankiki05@gmail.com>
This commit was merged in pull request #106.
This commit is contained in:
@@ -5,6 +5,8 @@ import (
|
||||
"errors"
|
||||
"log"
|
||||
"net/url"
|
||||
"sort"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"bookmarkmanager/backend/internal/store"
|
||||
@@ -28,15 +30,11 @@ type BrowserCoverFetcher interface {
|
||||
// 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.
|
||||
// - Cooldowns are how long a series rests since its own last check. Browser-
|
||||
// backed sites use the longer BrowserCooldown.
|
||||
//
|
||||
// Cooldowns are 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.
|
||||
// Every Site gets its own Poll Lane: one independent stream of Polls with its
|
||||
// own pace, running concurrently with every other Site's (issue #100). Rest
|
||||
// time and gap live in the Site registry, not here — see sites.go. Rest is
|
||||
// enforced by the WHERE clause in DueForLatestCheck rather than by any timer;
|
||||
// the gap is enforced by the Lane sleeping between fetches.
|
||||
type Poller struct {
|
||||
Store *store.Store
|
||||
Fetch Fetcher
|
||||
@@ -50,11 +48,19 @@ type Poller struct {
|
||||
// same failure-isolated prefetch path.
|
||||
CoverBytesFetch CoverBytesFetcher
|
||||
Now func() time.Time // injected so tests can freeze it
|
||||
Cooldown time.Duration
|
||||
BrowserCooldown time.Duration
|
||||
Interval time.Duration
|
||||
Stagger time.Duration
|
||||
Batch int
|
||||
|
||||
// refuseUntil gates a Site's Lane after it refused twice in one run: no
|
||||
// Series of that Site is attempted again before this time (issue #100).
|
||||
// browserDownAt is when a browser Lane last lost the sidecar; the other
|
||||
// browser Lanes skip their passes for the next refuseBackoff, so a
|
||||
// restarting Chrome does not stamp one Series per pass per Lane (story 20).
|
||||
mu sync.Mutex
|
||||
refuseUntil map[string]time.Time
|
||||
browserDownAt time.Time
|
||||
// coverWG tracks in-flight cover work. Covers heal in the background so a
|
||||
// slow cover host cannot delay the next Series-page Poll; tests join it
|
||||
// before asserting on cover fetches.
|
||||
coverWG sync.WaitGroup
|
||||
}
|
||||
|
||||
// fillBlankCover gives a Series its Cover when it has none. The blank state is
|
||||
@@ -76,7 +82,14 @@ func (p *Poller) fillBlankCover(ctx context.Context, sr store.Series, cover stri
|
||||
if cover == "" {
|
||||
return
|
||||
}
|
||||
p.storeCover(ctx, sr, cover)
|
||||
// Like healCover, the fill runs in the background: a large import of
|
||||
// blanks would otherwise pay one og:image fetch per Series against the
|
||||
// Lane's gap (issue #100, story 12).
|
||||
p.coverWG.Add(1)
|
||||
go func() {
|
||||
defer p.coverWG.Done()
|
||||
p.storeCover(ctx, sr, cover)
|
||||
}()
|
||||
}
|
||||
|
||||
// prefetchCover heals Series that already carry a third-party source URL but
|
||||
@@ -143,72 +156,248 @@ func fetcherFor(site string, browser, tls Fetcher) Fetcher {
|
||||
return nil
|
||||
}
|
||||
|
||||
// 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.
|
||||
// Run polls until ctx is cancelled: one goroutine per Site Lane, each pacing
|
||||
// itself by the Site's effective gap. Lanes share nothing but the store and
|
||||
// the browser fetcher's single tab (BrowserFetcher serializes itself), so one
|
||||
// hostile Site burns only its own budget.
|
||||
// laneNames returns every registry Site in the deterministic order both Run
|
||||
// and runOnce iterate: sorted, so lane behaviour and its tests agree on who
|
||||
// runs first.
|
||||
func laneNames() []string {
|
||||
names := make([]string, 0, len(sites))
|
||||
for name := range sites {
|
||||
names = append(names, name)
|
||||
}
|
||||
sort.Strings(names)
|
||||
return names
|
||||
}
|
||||
|
||||
func (p *Poller) Run(ctx context.Context) {
|
||||
log.Printf("latest-chapter poller: interval=%s cooldown=%s browser-cooldown=%s batch=%d stagger=%s",
|
||||
p.Interval, p.Cooldown, p.BrowserCooldown, p.Batch, p.Stagger)
|
||||
t := time.NewTicker(p.Interval)
|
||||
defer t.Stop()
|
||||
names := laneNames()
|
||||
log.Printf("latest-chapter poller: %d lanes, rest=%s gap=%s", len(names), defaultRest, defaultGap)
|
||||
for _, name := range names {
|
||||
go p.lane(ctx, name)
|
||||
}
|
||||
<-ctx.Done()
|
||||
log.Println("latest-chapter poller: stopped")
|
||||
}
|
||||
|
||||
// lane is one Site's Poll Lane: one pass, then sleep the pace the pass
|
||||
// reported, then another pass, until ctx is cancelled. The sleep is the whole
|
||||
// pace discipline — a pass that fetched nothing still reports its gap so the
|
||||
// Lane wakes often enough to notice Series as they become due. The pass shares
|
||||
// the Poller's browser-down state, so a sidecar loss is noticed once and the
|
||||
// other browser Lanes skip passes until the backoff window decays.
|
||||
func (p *Poller) lane(ctx context.Context, name string) {
|
||||
for {
|
||||
pace := p.runLanePass(ctx, name, true)
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
log.Println("latest-chapter poller: stopped")
|
||||
return
|
||||
case <-t.C:
|
||||
p.runOnce(ctx)
|
||||
case <-time.After(pace):
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// runOnce processes one batch of due series.
|
||||
// runOnce processes one round: one pass of every Lane, back to back, no real
|
||||
// time passing. This is the deterministic entry point the test suite drives a
|
||||
// round at a time. The production Run loop does the same work paced by its own
|
||||
// sleeps; pacing is the only difference.
|
||||
func (p *Poller) runOnce(ctx context.Context) {
|
||||
for _, name := range laneNames() {
|
||||
p.runLanePass(ctx, name, false)
|
||||
}
|
||||
}
|
||||
|
||||
// runLanePass processes one pass of one Site's Lane: select the due Series,
|
||||
// pace through them, and report how long the Lane should wait before its next
|
||||
// pass. paced spaces consecutive fetches by the Site's effective gap — the
|
||||
// production Lane's rate limit; the deterministic test entry runs back to back.
|
||||
func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.Duration {
|
||||
now := p.Now()
|
||||
cutoff := now.Add(-p.Cooldown).UnixMilli()
|
||||
browserCutoff := now.Add(-p.BrowserCooldown).UnixMilli()
|
||||
due, err := p.Store.DueForLatestCheck(cutoff, browserCutoff, browserBackedSites(), p.Batch)
|
||||
if err != nil {
|
||||
log.Printf("latest poll: due query: %v", err)
|
||||
return
|
||||
if until := p.refusalBackoff(name); now.Before(until) {
|
||||
// Cooling down after a refusal: do not attempt this Site at all.
|
||||
return until.Sub(now)
|
||||
}
|
||||
if isBrowserSite(name) {
|
||||
if downFor, down := p.browserDownFor(now); down && downFor < refuseBackoff {
|
||||
// A sibling browser Lane lost the sidecar within the backoff
|
||||
// window: skip this pass, so a restarting Chrome does not stamp
|
||||
// this Site's Series one pass at a time. After refuseBackoff the
|
||||
// flag decays and the Lane probes again (issue #100, story 20).
|
||||
log.Printf("latest poll %s: browser lane skipping pass (sidecar down %s ago)", name, downFor)
|
||||
return refuseBackoff - downFor
|
||||
}
|
||||
}
|
||||
s := sites[name]
|
||||
f := fetcherFor(name, p.BrowserFetch, p.Fetch)
|
||||
if f == nil {
|
||||
// No fetcher at all right now (browser absent, no fallback): every
|
||||
// Series stays unstamped and due, so a browser that appears after a
|
||||
// restart finds its full queue waiting (issue #100).
|
||||
return defaultGap
|
||||
}
|
||||
|
||||
due, err := p.Store.DueForLatestCheck(name, now.Add(-s.Rest).UnixMilli())
|
||||
if err != nil {
|
||||
log.Printf("latest poll %s: due query: %v", name, err)
|
||||
return defaultGap
|
||||
}
|
||||
if s.Browser != nil && f == p.BrowserFetch && !browserWakeDue(due, now, s.Rest) {
|
||||
// Below both thresholds Chrome stays asleep (ADR-0005 on-demand
|
||||
// browser): waking it for a single Poll would cost a challenge solve
|
||||
// per request.
|
||||
return defaultGap
|
||||
}
|
||||
if s.Browser != nil {
|
||||
// Browser Lanes share one tab, so their combined ceiling is about 360
|
||||
// Polls an hour. When they cannot keep up, the wait past the rest time
|
||||
// grows — log by how much, every pass, so the decision to give them
|
||||
// more pages is made from a measurement rather than a guess.
|
||||
if behind := maxSeriesWait(due, now, s.Rest) - s.Rest; behind > 0 {
|
||||
log.Printf("latest poll %s: browser lane behind by %s (browser Sites cannot keep up with the hour)", name, behind)
|
||||
}
|
||||
}
|
||||
|
||||
eligible, err := p.Store.EligibleSeriesCount(name)
|
||||
if err != nil {
|
||||
log.Printf("latest poll %s: eligible count: %v", name, err)
|
||||
return defaultGap
|
||||
}
|
||||
gap, clamped := effectiveGap(s, eligible)
|
||||
if clamped {
|
||||
log.Printf("latest poll %s: gap clamped to %s floor (eligible series=%d)", name, minGap, eligible)
|
||||
}
|
||||
if eligible == 0 {
|
||||
// Nothing to poll for the foreseeable future; sleep a full rest instead
|
||||
// of re-querying every gap.
|
||||
return s.Rest
|
||||
}
|
||||
|
||||
refusals := 0
|
||||
checked := 0
|
||||
for i, sr := 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).
|
||||
stopped := false
|
||||
if i > 0 && p.Stagger > 0 {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
stopped = true
|
||||
case <-time.After(p.Stagger):
|
||||
}
|
||||
}
|
||||
if stopped {
|
||||
if refusals >= 2 {
|
||||
// This Site refused twice in a row: the remaining Series are left
|
||||
// unstamped and due, and the Lane waits refuseBackoff before
|
||||
// trying it again.
|
||||
break
|
||||
}
|
||||
p.checkOne(ctx, sr)
|
||||
if paced && i > 0 {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
break
|
||||
case <-time.After(gap):
|
||||
}
|
||||
if ctx.Err() != nil {
|
||||
break
|
||||
}
|
||||
}
|
||||
if err := p.checkOne(ctx, sr); err != nil {
|
||||
switch {
|
||||
case errors.Is(err, errChallengeHeld):
|
||||
refusals++
|
||||
case errors.Is(err, errBrowserInterrupted):
|
||||
p.setBrowserDown(now)
|
||||
log.Printf("latest poll %s: browser unreachable, browser lanes skipping passes for %s", name, refuseBackoff)
|
||||
return gap
|
||||
default:
|
||||
refusals = 0
|
||||
}
|
||||
} else {
|
||||
refusals = 0
|
||||
}
|
||||
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)
|
||||
if checked > 0 {
|
||||
log.Printf("latest poll %s: due=%d checked=%d", name, len(due), checked)
|
||||
}
|
||||
if refusals >= 2 {
|
||||
p.setRefusalBackoff(name, now.Add(refuseBackoff))
|
||||
log.Printf("latest poll %s: refused twice this run, waiting %s", name, refuseBackoff)
|
||||
return refuseBackoff
|
||||
}
|
||||
return gap
|
||||
}
|
||||
|
||||
func (p *Poller) refusalBackoff(name string) time.Time {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
return p.refuseUntil[name]
|
||||
}
|
||||
|
||||
func (p *Poller) setRefusalBackoff(name string, until time.Time) {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
if p.refuseUntil == nil {
|
||||
p.refuseUntil = make(map[string]time.Time)
|
||||
}
|
||||
p.refuseUntil[name] = until
|
||||
}
|
||||
|
||||
// setBrowserDown records when a browser Lane lost the sidecar. It is Poller
|
||||
// state rather than pass state so the other browser Lanes see it too.
|
||||
func (p *Poller) setBrowserDown(now time.Time) {
|
||||
p.mu.Lock()
|
||||
p.browserDownAt = now
|
||||
p.mu.Unlock()
|
||||
}
|
||||
|
||||
// browserDownFor reports how long the sidecar has been down and that it is
|
||||
// down at all — the zero time means never down, which must not read as a
|
||||
// zero-duration loss. The window decays: once refuseBackoff passes without a
|
||||
// fresh loss, Lanes probe again.
|
||||
func (p *Poller) browserDownFor(now time.Time) (time.Duration, bool) {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
if p.browserDownAt.IsZero() {
|
||||
return 0, false
|
||||
}
|
||||
return now.Sub(p.browserDownAt), true
|
||||
}
|
||||
|
||||
// isBrowserSite reports whether the registry routes this Site's page through
|
||||
// the browser sidecar.
|
||||
func isBrowserSite(name string) bool {
|
||||
return sites[name].Browser != nil
|
||||
}
|
||||
|
||||
// browserWakeDue reports whether a browser Lane may start a run: five or more
|
||||
// of its Series are due, or any one of them has been due for browserWakeAge.
|
||||
// Below both thresholds the Lane leaves Chrome asleep — Series Polled together
|
||||
// become due together, so the group naturally stays clustered, and the age
|
||||
// rule exists to stop a Series that drifted out of the group from starving.
|
||||
func browserWakeDue(due []store.Series, now time.Time, rest time.Duration) bool {
|
||||
if len(due) >= browserWakeCount {
|
||||
return true
|
||||
}
|
||||
return maxSeriesWait(due, now, rest) >= browserWakeAge
|
||||
}
|
||||
|
||||
// maxSeriesWait returns how long the most-overdue of the due Series has been
|
||||
// waiting past its due moment (0 when due is empty).
|
||||
func maxSeriesWait(due []store.Series, now time.Time, rest time.Duration) time.Duration {
|
||||
var oldest time.Duration
|
||||
for _, sr := range due {
|
||||
if w := now.Sub(time.UnixMilli(sr.LatestCheckedAt).Add(rest)); w > oldest {
|
||||
oldest = w
|
||||
}
|
||||
}
|
||||
return oldest
|
||||
}
|
||||
|
||||
// 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 *Poller) checkOne(ctx context.Context, sr store.Series) {
|
||||
// Lane or take down the process. The returned error is the page read's
|
||||
// classified outcome so the Lane can tell a refusal from a loss of the
|
||||
// browser; non-classified failures still return nil-equivalent behaviour.
|
||||
func (p *Poller) checkOne(ctx context.Context, sr store.Series) error {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
log.Printf("latest poll %q: recovered from panic: %v", sr.Key(), r)
|
||||
@@ -216,44 +405,46 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) {
|
||||
}()
|
||||
|
||||
// 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).
|
||||
// mid-request still consumes the rest. Otherwise a renamed or deleted
|
||||
// series would be retried on every single pass forever. The userscript
|
||||
// stamps in the same order and for the same reason (L471-473). A Series
|
||||
// never reaches checkOne without a fetcher — runLanePass skips those — so
|
||||
// the stamp means "attempted", and an untried Series stays due.
|
||||
if err := p.Store.MarkLatestChecked(sr.Site, sr.SeriesID, p.Now().UnixMilli()); err != nil {
|
||||
log.Printf("latest poll %q: mark checked: %v", sr.Key(), err)
|
||||
return
|
||||
return nil
|
||||
}
|
||||
|
||||
facts, err := readSeriesPage(ctx, sr.Site, sr.SeriesURL, p.BrowserFetch, p.Fetch)
|
||||
if err != nil {
|
||||
switch {
|
||||
case errors.Is(err, errNotFetchable):
|
||||
// The cooldown above is already consumed, so a row that never
|
||||
// passes the gate is retried at cooldown pace rather than
|
||||
// The rest above is already consumed, so a row that never
|
||||
// passes the gate is retried at rest pace rather than
|
||||
// hot-looping.
|
||||
log.Printf("latest poll %q: not fetchable: site=%q url=%q", sr.Key(), sr.Site, sr.SeriesURL)
|
||||
return
|
||||
return err
|
||||
case errors.Is(err, errNoFetcher):
|
||||
log.Printf("latest poll %q: no fetcher for site %q", sr.Key(), sr.Site)
|
||||
return
|
||||
return err
|
||||
}
|
||||
// A legacy cover heals independently of the page read: its source may
|
||||
// answer — a CDN — while the origin does not, so a fetch failure does
|
||||
// not skip the heal, matching the order the shared read replaced.
|
||||
p.prefetchCover(ctx, sr)
|
||||
p.healCover(ctx, sr)
|
||||
log.Printf("latest poll %q: %v", sr.Key(), err)
|
||||
return
|
||||
return err
|
||||
}
|
||||
// A legacy cover source is healed independently of the page read.
|
||||
p.prefetchCover(ctx, sr)
|
||||
p.healCover(ctx, sr)
|
||||
// Cover fill is independent of the chapter signal: a page that lost its
|
||||
// chapter list may keep its og:image, and a blank Series heals either way.
|
||||
p.fillBlankCover(ctx, sr, facts.Cover)
|
||||
if !facts.HasLatest {
|
||||
// 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.
|
||||
// already stamped, so this waits out a rest instead of hot-looping.
|
||||
log.Printf("latest poll %q: no chapter links in %d bytes", sr.Key(), facts.BodyLen)
|
||||
return
|
||||
return nil
|
||||
}
|
||||
|
||||
// Equality, not >, mirroring the userscript (L427): a site that retracts a
|
||||
@@ -261,7 +452,7 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) {
|
||||
// against the due-query snapshot; a concurrent write in between only costs
|
||||
// one redundant UPDATE of the same absolute value, never a wrong one.
|
||||
if sr.LatestChapterNum != nil && *sr.LatestChapterNum == facts.Latest.Num {
|
||||
return
|
||||
return nil
|
||||
}
|
||||
|
||||
// Series-level write: the row is shared, so one update refreshes every
|
||||
@@ -270,9 +461,29 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) {
|
||||
// the list.
|
||||
if err := p.Store.SetLatestChapter(sr.Site, sr.SeriesID, facts.Latest.Label, facts.Latest.Num); err != nil {
|
||||
log.Printf("latest poll %q: set latest chapter: %v", sr.Key(), err)
|
||||
return
|
||||
return nil
|
||||
}
|
||||
log.Printf("latest poll %q: latest is now %s", sr.Key(), facts.Latest.Label)
|
||||
return nil
|
||||
}
|
||||
|
||||
// healCover runs prefetchCover in the background. Cover bytes come from a
|
||||
// different host — often a CDN — and heal once in a Series's life, so they
|
||||
// must not consume a Lane's gap: a large import with many blanks would
|
||||
// otherwise make every Latest Chapter go stale behind a slow image host
|
||||
// (issue #100).
|
||||
func (p *Poller) healCover(ctx context.Context, sr store.Series) {
|
||||
p.coverWG.Add(1)
|
||||
go func() {
|
||||
defer p.coverWG.Done()
|
||||
p.prefetchCover(ctx, sr)
|
||||
}()
|
||||
}
|
||||
|
||||
// waitCovers blocks until every in-flight cover heal finishes. Tests call it
|
||||
// after a round before asserting on cover fetches.
|
||||
func (p *Poller) waitCovers() {
|
||||
p.coverWG.Wait()
|
||||
}
|
||||
|
||||
// fetchableSeriesURL reports whether site is a Site the registry knows and
|
||||
|
||||
Reference in New Issue
Block a user