Files
mangaBookmark/backend/internal/latest/poller.go
T
sulthan 56afb9f237 feat: a Reader's Sighting defers a Poll of a solitary Series (#103)
A userscript PUT already carries the Latest Chapter the Reader's own browser
read; it now stands in for a Poll where being wrong can hurt nobody else.
Deferral lives in the due query beside the rest cutoff: one Bookmark, sighted
within one rest, under the six-rest ceiling. checkOne judges the Reader whose
report raised the value off the comparison it already makes - three
contradictions stop them deferring, twenty confirmations forgive.

ADR-0011 records the trust model and the rejected alternatives.
2026-08-16 15:17:30 +07:00

559 lines
22 KiB
Go

package latest
import (
"context"
"errors"
"log"
"net/url"
"sort"
"sync"
"time"
"bookmarkmanager/backend/internal/store"
)
// 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)
}
// BrowserCoverFetcher retrieves one cover's bytes through the browser-backed
// path — the only route that clears the challenge kagane's and comix's image
// URLs answer a plain fetch with. Satisfied by BrowserFetcher.
type BrowserCoverFetcher interface {
Image(ctx context.Context, imageURL string) (body []byte, contentType string, err error)
}
// Poller 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.
//
// 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
// BrowserFetch handles sites behind a JavaScript challenge that Fetch
// cannot clear. Nil disables those sites entirely rather than falling back
// to Fetch, which would only ever retrieve a challenge page.
BrowserFetch Fetcher
// CoverFetch is optional; failures are logged and never affect the chapter poll.
CoverFetch BrowserCoverFetcher
// CoverBytesFetch is optional; it handles plain-TLS sources through the
// same failure-isolated prefetch path.
CoverBytesFetch CoverBytesFetcher
Now func() time.Time // injected so tests can freeze it
// 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
// laneStates is the owner's page snapshot of each Lane's last pass
// (issue #102), keyed by Site. Guarded by mu; a Site appears only after
// its first pass, so a restart renders "no data yet" rather than zeroes.
laneStates map[string]LaneState
// 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
// what "no Cover yet" means on the wire (ADR-0007): permanently-blank rows
// created before acquisition existed, and rows whose creation-time fetch
// failed, both heal here. A non-blank CoverAddress is left alone — refetching
// would add a request per Series per cycle and change artwork under the Reader
// for no visible reason. A row that already carries a source URL is owned by
// prefetchCover instead; this path only records a Cover address already
// extracted from the series page.
//
// Failures are logged against the Series and never returned: the chapter poll
// must not notice. A failed fill is retried the next time this Series is due;
// there is no separate retry queue.
func (p *Poller) fillBlankCover(ctx context.Context, sr store.Series, cover string) {
if sr.CoverAddress != "" || sr.Cover != "" {
return
}
if cover == "" {
return
}
// 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
// no stored address — the state left by client-supplied covers before
// acquisition moved server-side. Every Site takes the same path; fetchCoverBytes
// routes by URL shape, so browser-claimed URLs still need the sidecar. New
// blanks have no source URL and go through fillBlankCover from the series page
// instead.
func (p *Poller) prefetchCover(ctx context.Context, sr store.Series) {
if sr.Cover == "" || sr.CoverAddress != "" {
return
}
body, contentType, found, err := p.Store.GetCover(sr.Cover)
if err != nil {
log.Printf("latest poll %q: read cover: %v", sr.Key(), err)
return
}
if found {
if err := p.Store.SetSeriesCover(sr.Site, sr.SeriesID, sr.Cover, body, contentType); err != nil {
log.Printf("latest poll %q: persist cover: %v", sr.Key(), err)
}
return
}
p.storeCover(ctx, sr, sr.Cover)
}
// storeCover fetches bytes for sourceURL and points the Series at them. Every
// failure is logged against the Series and swallowed so the chapter poll
// cannot see it.
func (p *Poller) storeCover(ctx context.Context, sr store.Series, sourceURL string) {
bytes, contentType, err := fetchCoverBytes(ctx, sourceURL, p.CoverFetch, p.CoverBytesFetch)
if err != nil {
log.Printf("latest poll %q: fetch cover %s: %v", sr.Key(), sourceURL, err)
return
}
if err := p.Store.SetSeriesCover(sr.Site, sr.SeriesID, sourceURL, bytes, contentType); err != nil {
log.Printf("latest poll %q: persist cover: %v", sr.Key(), err)
}
}
// fetcherFor returns the fetcher a site's page needs, or nil when the site
// cannot be fetched at all right now. A Site whose registry entry carries a
// Browser read — kagane, comix and novelfull, all behind a Cloudflare
// JavaScript challenge no TLS fingerprint clears — prefers the browser; when it
// is absent, the entry's Fallback decides whether plain TLS may take over. One
// routing rule for the poll and the acquirer, so the two cannot drift apart.
func fetcherFor(site string, browser, tls Fetcher) Fetcher {
s, known := sites[site]
if !known {
// No registry entry means nothing to fetch or parse; fail closed even
// though the only caller gates first, so a future caller that skips
// the gate cannot hand an arbitrary https URL to the TLS fetcher.
return nil
}
if s.Browser == nil {
return tls
}
if browser != nil {
return browser
}
if s.Browser.Fallback {
return tls
}
return nil
}
// 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) {
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():
return
case <-time.After(pace):
}
}
}
// 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()
// Snapshot this pass for the owner's page (issue #102). Recorded on every
// return path, with the figures filled in where the pass computes them.
st := LaneState{Site: name, LastRun: now, Browser: isBrowserSite(name)}
defer func() { p.recordLaneState(st) }()
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).
st.Gap = defaultGap
return defaultGap
}
due, err := p.Store.DueForLatestCheck(name,
now.Add(-s.Rest).UnixMilli(), now.Add(-sightingCeiling).UnixMilli())
if err != nil {
log.Printf("latest poll %s: due query: %v", name, err)
st.Gap = defaultGap
return defaultGap
}
st.Due = len(due)
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. The Lane still paces at the default gap, which is what
// the owner's page must show rather than a zero.
st.Gap = defaultGap
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)
st.Gap = defaultGap
return defaultGap
}
gap, clamped := effectiveGap(s, eligible)
st.Gap, st.Clamped = gap, clamped
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
for i, sr := range due {
if ctx.Err() != nil {
break
}
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
}
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
}
st.Checked++
}
if st.Checked > 0 {
log.Printf("latest poll %s: due=%d checked=%d", name, len(due), st.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
// 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)
}
}()
// Stamped before the fetch, not after, so an error, a timeout, or a shutdown
// 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 nil
}
facts, err := readSeriesPage(ctx, sr.Site, sr.SeriesURL, p.BrowserFetch, p.Fetch)
if err != nil {
switch {
case errors.Is(err, errNotFetchable):
// 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 err
case errors.Is(err, errNoFetcher):
log.Printf("latest poll %q: no fetcher for site %q", sr.Key(), sr.Site)
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.healCover(ctx, sr)
log.Printf("latest poll %q: %v", sr.Key(), err)
return err
}
// A legacy cover source is healed independently of the page read.
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 rest instead of hot-looping.
log.Printf("latest poll %q: no chapter links in %d bytes", sr.Key(), facts.BodyLen)
return nil
}
// The Poll is the oracle for whatever Sighting last raised this Series
// (issue #103), and the judgement is free: the comparison below already
// exists, and no extra request is made to reach it.
p.judgeSighting(sr, facts.Latest.Num)
// Equality, not >, mirroring the userscript (L427): a site that retracts a
// chapter should correct the stored number downward. The comparison is
// 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 nil
}
// Series-level write: the row is shared, so one update refreshes every
// bookmark joining to it, and the bookmark's updated_at is never touched —
// a newly published chapter is not reading progress and must not reorder
// 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 nil
}
log.Printf("latest poll %q: latest is now %s", sr.Key(), facts.Latest.Label)
return nil
}
// judgeSighting settles the Sighting the Series' stored Latest Chapter is owed
// to, if any, against what the Site actually publishes. The asymmetry is the
// whole of the detection rule and is what keeps it free of false alarms: a Poll
// finding a *lower* number than stored means the Reader who raised it reported
// a chapter that does not exist, while a Poll finding a higher one is only the
// Site publishing since and means nothing about the report. Equality confirms
// the report, which is how an honest Reader earns back a mark.
//
// A Series with no attribution — the stored value is a Poll's own, or a
// previous Poll already judged the report — is nobody's to answer for.
func (p *Poller) judgeSighting(sr store.Series, found float64) {
if sr.LatestRaisedBy == nil || sr.LatestChapterNum == nil {
return
}
stored := *sr.LatestChapterNum
if found > stored {
return
}
if found < stored {
// Logged with both numbers and the Reader, because that is what tells a
// broken Site adapter (which marks every Reader of that Site at once)
// from one Reader deliberately lying.
log.Printf("latest poll %q: sighting contradicted: reader %d raised it to %v, site publishes %v",
sr.Key(), *sr.LatestRaisedBy, stored, found)
}
if err := p.Store.RecordSightingOutcome(sr.Site, sr.SeriesID, *sr.LatestRaisedBy, found == stored); err != nil {
log.Printf("latest poll %q: record sighting outcome: %v", sr.Key(), err)
}
}
// 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
// seriesURL is safe to hand to a fetcher: an https URL whose host matches the
// Site's pinned hostname exactly. series_url comes from client-supplied PUT
// bodies, so this is a defence against the poller being used to probe
// arbitrary hosts from the server's own network position, not just a check
// against wasted requests. The pin guards different things per Site — a
// browser Site guards a control that executes JavaScript and carries cookies,
// a parser Site guards a wasted request — but the rule is one rule, from the
// registry.
func fetchableSeriesURL(site, seriesURL string) bool {
s, known := sites[site]
if !known {
return false
}
u, err := url.Parse(seriesURL)
if err != nil {
return false
}
return u.Scheme == "https" && u.Hostname() == s.Host
}