7b22460f5e
Rollout note: the owner_notices table starts empty, so the first pass after deploy sends for conditions already true — correct per one-row-per-episode; say so rather than have it reported as a bug. Prod step: create the webhook, set DISCORD_WEBHOOK_URL on the deployment, redeploy — unset is silent by design, and without that step the feature ships dark. Security invariants preserved: the webhook address is a secret in the class of TOKEN_KEY (never logged, never on a config-printing line), and the owner gate is unchanged.
915 lines
36 KiB
Go
915 lines
36 KiB
Go
package latest
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"log"
|
|
"net/url"
|
|
"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
|
|
// Notify delivers owner notices. Nil disables the whole path (issue #171):
|
|
// the poller is not the place a missing webhook becomes an error.
|
|
Notify Notifier
|
|
Now func() time.Time // injected so tests can freeze it
|
|
// eligibleCount reports how many of a Site's Series are eligible for
|
|
// polling, defaulting to Store.EligibleSeriesCount. Injected so tests can
|
|
// fail the count alone: the eligible query shares the due query's tables,
|
|
// so no real store failure can reach this path without breaking the due
|
|
// query first (issue #141).
|
|
eligibleCount func(site string) (int, error)
|
|
|
|
// 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).
|
|
// The stamp is durable — the pass gate reads it from the store, so a
|
|
// restart does not forget the refusal; nothing of it lives in memory.
|
|
// 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
|
|
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
|
|
// 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)
|
|
}()
|
|
}
|
|
|
|
// replaceCover is the Forced Poll's Cover path: the owner asked to accept the
|
|
// page as it now stands, so where fillBlankCover leaves a non-blank Cover
|
|
// alone (ADR-0007) this writes through whatever the page's Cover URL answers
|
|
// with, whether one exists or not. The accepted consequence (issue #135):
|
|
// refreshing the Cover and re-reading the chapters are one act — there is no
|
|
// Cover-only refetch.
|
|
func (p *Poller) replaceCover(ctx context.Context, sr store.Series, cover string) {
|
|
if cover == "" {
|
|
return
|
|
}
|
|
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: a
|
|
// fill-only write for an ordinary pass, a write-through for a forced one
|
|
// (issue #135). 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 !sr.Forced {
|
|
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)
|
|
}
|
|
return
|
|
}
|
|
// The forced write replaces whether or not a Cover exists, and the row
|
|
// then tells the three outcomes apart: a blank filled, identical artwork
|
|
// re-served — an honest no-op — or a replacement whose previous address
|
|
// is stranded and reclaimed below. A failed reclaim is logged and the
|
|
// stranded bytes stay served until a later call reclaims them.
|
|
previous, current, err := p.Store.ReplaceSeriesCover(sr.Site, sr.SeriesID, sourceURL, bytes, contentType)
|
|
if err != nil {
|
|
log.Printf("latest poll %q: persist cover: %v", sr.Key(), err)
|
|
return
|
|
}
|
|
switch {
|
|
case previous == "":
|
|
log.Printf("latest poll %q: cover filled at %s", sr.Key(), current)
|
|
case previous == current:
|
|
log.Printf("latest poll %q: cover unchanged, the site re-serves the same bytes", sr.Key())
|
|
default:
|
|
if err := p.Store.ReclaimCover(previous); err != nil {
|
|
log.Printf("latest poll %q: reclaim cover %s: %v", sr.Key(), previous, err)
|
|
}
|
|
log.Printf("latest poll %q: cover replaced %s -> %s", sr.Key(), previous, current)
|
|
}
|
|
}
|
|
|
|
// 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 { return SiteNames() }
|
|
|
|
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)
|
|
}
|
|
}
|
|
|
|
// One skip value per way a Lane Pass can return before its loop (issue #141);
|
|
// empty means the pass reached the loop. The values are wire strings — stored
|
|
// in poll_passes and read by the Lanes page — so they are stable, not prose.
|
|
const (
|
|
// Exported so the web layer renders a skip's reason without retyping the
|
|
// wire string (issue #145); the values are storage and page-stable.
|
|
SkipPaused = "paused" // the pause row was read at the top
|
|
SkipRefusing = "refusing" // refusal backoff
|
|
SkipSidecarDown = "sidecar-down" // a sibling browser Lane lost Chrome
|
|
SkipNoFetcher = "no-fetcher" // browser Site, no browser configured, no fallback
|
|
SkipDueQuery = "due-query" // the due query failed
|
|
SkipAsleep = "asleep" // under both browser wake thresholds
|
|
SkipEligibleCount = "eligible-count" // the eligible count failed
|
|
SkipNothingEligible = "nothing-eligible" // nothing eligible; sleeps a full rest
|
|
)
|
|
|
|
// readOutcome classifies one Series read for the pass row's outcome counts
|
|
// (issue #141). The classification the read already makes is counted, never a
|
|
// second taxonomy: refused is the Site holding a challenge, unreachable the
|
|
// browser interrupting, noChapter a 200 with real HTML but no chapter links,
|
|
// unfetchable the host pin or a missing fetcher, notFound a 4xx other than
|
|
// the 403 refusal — the Site answered with a client status — and errors
|
|
// everything else.
|
|
type readOutcome int
|
|
|
|
const (
|
|
outcomeSuccess readOutcome = iota
|
|
outcomeRefused
|
|
outcomeUnreachable
|
|
outcomeNoChapter
|
|
outcomeUnfetchable
|
|
outcomeError
|
|
outcomeNotFound
|
|
)
|
|
// word returns the wire spelling this outcome stores in poll_failures — the
|
|
// same strings the pass log's columns use (C1, issue #164). Success, refusal
|
|
// and browser loss return "" so recordFailure's "no statement" case is one
|
|
// return: a challenge or a lost sidecar is no evidence about any particular
|
|
// Series (ADR-0016).
|
|
func (o readOutcome) word() string {
|
|
switch o {
|
|
case outcomeNotFound:
|
|
return "not_found"
|
|
case outcomeNoChapter:
|
|
return "no_chapter"
|
|
case outcomeUnfetchable:
|
|
return "unfetchable"
|
|
case outcomeError:
|
|
return "errors"
|
|
}
|
|
return ""
|
|
}
|
|
|
|
// outcomeCounts are the six named outcome counts of one pass. A success
|
|
// count is derived, never stored: checked minus the five named failures,
|
|
// with unreachable excluded because the sidecar-loss path returns before the
|
|
// checked counter increments (issue #141).
|
|
type outcomeCounts struct {
|
|
refused, unreachable, noChapter, unfetchable, errors, notFound int
|
|
}
|
|
|
|
func (c *outcomeCounts) add(o readOutcome) {
|
|
switch o {
|
|
case outcomeRefused:
|
|
c.refused++
|
|
case outcomeUnreachable:
|
|
c.unreachable++
|
|
case outcomeNoChapter:
|
|
c.noChapter++
|
|
case outcomeUnfetchable:
|
|
c.unfetchable++
|
|
case outcomeNotFound:
|
|
c.notFound++
|
|
case outcomeError:
|
|
c.errors++
|
|
}
|
|
}
|
|
|
|
// passRecord is what one pass's durable row will be: the skip value and
|
|
// outcome counts filled in along the pass's return path. recordPass assembles
|
|
// the row, so every exit records exactly once.
|
|
type passRecord struct {
|
|
site string
|
|
ranAt int64
|
|
skip string
|
|
counts outcomeCounts
|
|
}
|
|
|
|
// lanePassRetention is how far back a Lane's pass log is kept. It is not the
|
|
// display window: retention is how far back a question can reach, and the
|
|
// window is what the owner is shown (issue #139).
|
|
const lanePassRetention = 14 * 24 * time.Hour
|
|
|
|
// 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()
|
|
// Durable pass log (issue #141): one row per exit. The figures are filled
|
|
// in as the pass measures them; a pass that returns before measuring
|
|
// carries the previous pass's forward inside recordPass.
|
|
fig := passFigures{}
|
|
rec := passRecord{site: name, ranAt: now.UnixMilli()}
|
|
defer func() { p.recordPass(ctx, rec, fig) }()
|
|
|
|
// One Lane row read at the top of a pass, serving two gates (issue #139).
|
|
// Both stamps outlive our process, so the gates read the durable row
|
|
// rather than memory: a refusal is the Site's mood and a pause the
|
|
// owner's order, and neither is lost to a restart.
|
|
pausedUntil, refuseUntil, err := p.Store.LaneGates(name)
|
|
if err != nil {
|
|
// Fail open: a store that cannot answer the gate cannot record the
|
|
// pass either, and one Lane must not stall on its own gate read.
|
|
log.Printf("latest poll %s: lane gates: %v", name, err)
|
|
}
|
|
if pausedUntil > now.UnixMilli() {
|
|
// Paused ahead of the refusal check: no Series is touched, so the
|
|
// queue stays intact for when the pause lifts (issue #141, #147).
|
|
rec.skip = SkipPaused
|
|
log.Printf("latest poll %s: paused until %s, skipping pass", name, time.UnixMilli(pausedUntil).Format(time.RFC3339))
|
|
return time.Duration(pausedUntil-now.UnixMilli()) * time.Millisecond
|
|
}
|
|
if refuseUntil > now.UnixMilli() {
|
|
// Cooling down after a refusal: do not attempt this Site at all.
|
|
rec.skip = SkipRefusing
|
|
return time.Duration(refuseUntil-now.UnixMilli()) * time.Millisecond
|
|
}
|
|
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).
|
|
rec.skip = SkipSidecarDown
|
|
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).
|
|
rec.skip = SkipNoFetcher
|
|
fig.Gap = defaultGap
|
|
return defaultGap
|
|
}
|
|
|
|
due, err := p.Store.DueForLatestCheck(name, now.Add(-s.Rest).UnixMilli(),
|
|
now.Add(-sightingCeilingRests*s.Rest).UnixMilli())
|
|
if err != nil {
|
|
rec.skip = SkipDueQuery
|
|
log.Printf("latest poll %s: due query: %v", name, err)
|
|
fig.Gap = defaultGap
|
|
return defaultGap
|
|
}
|
|
fig.Due = len(due)
|
|
if s.Browser != nil && f == p.BrowserFetch && !browserWakeDue(due, now, s.Rest) && !anyForced(due) {
|
|
// Below both thresholds Chrome stays asleep (ADR-0005 on-demand
|
|
// browser): waking it for a single Poll would cost a challenge solve
|
|
// per request. A forced Series is the one exception — a human asking
|
|
// is not the machine waking itself (issue #146). The Lane still paces
|
|
// at the default gap, which is what the owner's page must show rather
|
|
// than a zero.
|
|
rec.skip = SkipAsleep
|
|
fig.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.countEligible(name)
|
|
if err != nil {
|
|
rec.skip = SkipEligibleCount
|
|
log.Printf("latest poll %s: eligible count: %v", name, err)
|
|
fig.Gap = defaultGap
|
|
return defaultGap
|
|
}
|
|
gap, clamped := effectiveGap(s, eligible)
|
|
fig.Gap, fig.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.
|
|
rec.skip = SkipNothingEligible
|
|
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
|
|
}
|
|
}
|
|
outcome := p.checkOne(ctx, sr)
|
|
p.recordFailure(sr, outcome)
|
|
if outcome == outcomeUnreachable {
|
|
// The mid-loop browser loss writes an empty skip on purpose: the
|
|
// pass returns before the checked counter increments, so its row
|
|
// is stall-shaped (due > 0, checked 0, skip ''), and a stall is
|
|
// the exact signal this exit produces. A tenth skip value would
|
|
// make it legible but is deliberately not invented here.
|
|
rec.counts.add(outcome)
|
|
p.setBrowserDown(now)
|
|
log.Printf("latest poll %s: browser unreachable, browser lanes skipping passes for %s", name, RefuseBackoff)
|
|
return gap
|
|
}
|
|
if outcome == outcomeRefused {
|
|
refusals++
|
|
} else {
|
|
refusals = 0
|
|
}
|
|
rec.counts.add(outcome)
|
|
fig.Checked++
|
|
}
|
|
if fig.Checked > 0 {
|
|
log.Printf("latest poll %s: due=%d checked=%d", name, len(due), fig.Checked)
|
|
}
|
|
if refusals >= 2 {
|
|
// The refusal outlives the process: the durable stamp gates a restart,
|
|
// so a Site that just told us to back off is not re-probed.
|
|
if err := p.Store.SetLaneRefusal(name, now.Add(RefuseBackoff).UnixMilli()); err != nil {
|
|
log.Printf("latest poll %s: persist refusal: %v", name, err)
|
|
}
|
|
log.Printf("latest poll %s: refused twice this run, waiting %s", name, RefuseBackoff)
|
|
return RefuseBackoff
|
|
}
|
|
return gap
|
|
}
|
|
|
|
// passFigures are the numbers one pass measured for its durable row (issue
|
|
// #141): due and checked as the pass saw them, the pace it chose, and whether
|
|
// the gap sat on the floor. A pass that returned before measuring keeps the
|
|
// previous pass's figures via carry-forward in recordPass; the in-memory
|
|
// snapshot those once mirrored into is gone — the page reads the durable row
|
|
// now (issue #145).
|
|
type passFigures struct {
|
|
Due, Checked int
|
|
Gap time.Duration
|
|
Clamped bool
|
|
}
|
|
|
|
// recordPass writes the durable row for one pass (issue #141). Called deferred
|
|
// from runLanePass so every return path records exactly one row. A pass that
|
|
// never computed its own figures — its gap is zero — carries the previous
|
|
// pass's due, gap, clamped and checked forward rather than stating zeroes it
|
|
// did not measure; the skip column says why it declined, so the zeroes that
|
|
// remain (due-query, no-fetcher) read as explanations rather than
|
|
// measurements. Once the row is durable, the owner-notice judgement runs
|
|
// beside it (issue #171).
|
|
func (p *Poller) recordPass(ctx context.Context, rec passRecord, fig passFigures) {
|
|
row := store.LanePass{
|
|
Site: rec.site,
|
|
RanAt: rec.ranAt,
|
|
Skip: rec.skip,
|
|
Due: fig.Due,
|
|
Checked: fig.Checked,
|
|
GapMS: fig.Gap.Milliseconds(),
|
|
Clamped: fig.Clamped,
|
|
Refused: rec.counts.refused,
|
|
Unreachable: rec.counts.unreachable,
|
|
NoChapter: rec.counts.noChapter,
|
|
Unfetchable: rec.counts.unfetchable,
|
|
NotFound: rec.counts.notFound,
|
|
Errors: rec.counts.errors,
|
|
}
|
|
if row.GapMS == 0 {
|
|
// The pass never computed a gap, so it has no figures of its own:
|
|
// carry the previous pass's, in one latest-per-Site read — the
|
|
// recorder needs one Site, not six (issue #139).
|
|
if prev, ok, err := p.Store.LatestLanePass(rec.site); err != nil {
|
|
log.Printf("latest poll %s: previous pass: %v", rec.site, err)
|
|
} else if ok {
|
|
row.Due, row.Checked = prev.Due, prev.Checked
|
|
row.GapMS, row.Clamped = prev.GapMS, prev.Clamped
|
|
}
|
|
}
|
|
if err := p.Store.RecordLanePass(row, rec.ranAt-lanePassRetention.Milliseconds()); err != nil {
|
|
log.Printf("latest poll %s: record lane pass: %v", rec.site, err)
|
|
return
|
|
}
|
|
p.ownerNotices(ctx, row)
|
|
}
|
|
|
|
// ownerNotices judges the owner-notice conditions for the pass just recorded
|
|
// and fires (issue #171). It sits in recordPass because that deferred call is
|
|
// the one place every return path passes through: two of the four conditions
|
|
// occur on early returns and the success path can never see them. Per fault:
|
|
// NoticeSent → send → MarkNoticeSent, so a fault lasting a month sends one
|
|
// message, not one per pass; a condition absent from this pass's fault list
|
|
// forgets its episode, so the next occurrence sends again. Everything here is
|
|
// best-effort: a failed send, a failed store read and a failed notice write
|
|
// are all logged and never change the pass's outcome counts or its return
|
|
// value. The clear runs even when Notify is nil, so a deployment that turns
|
|
// the webhook off does not leave stale rows that suppress the first real
|
|
// notice after it is turned back on.
|
|
func (p *Poller) ownerNotices(ctx context.Context, row store.LanePass) {
|
|
faults := FaultsFrom(FaultInput{Passes: []store.LanePass{row}}, p.Now())
|
|
for _, f := range faults {
|
|
if p.Notify == nil {
|
|
continue
|
|
}
|
|
sent, err := p.Store.NoticeSent(f.Condition, f.Site)
|
|
if err != nil {
|
|
log.Printf("latest poll %s: notice sent: %v", row.Site, err)
|
|
continue
|
|
}
|
|
if sent {
|
|
continue
|
|
}
|
|
sentence, href := noticeFor(f, row, p.Now())
|
|
if err := p.Notify.Notify(ctx, f, sentence, href); err != nil {
|
|
// The stamp stays unset: no queue, no backoff — the condition is
|
|
// durable, so the next pass tries again while it holds.
|
|
log.Printf("latest poll %s: owner notice %s: %v", row.Site, f.Condition, err)
|
|
continue
|
|
}
|
|
if err := p.Store.MarkNoticeSent(f.Condition, f.Site, p.Now().UnixMilli()); err != nil {
|
|
log.Printf("latest poll %s: mark notice sent: %v", row.Site, err)
|
|
}
|
|
}
|
|
for _, cond := range ownerNoticeConditions {
|
|
if !hasFault(faults, cond, row.Site) {
|
|
if err := p.Store.ClearNotice(cond, row.Site); err != nil {
|
|
log.Printf("latest poll %s: clear owner notice: %v", row.Site, err)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// hasFault reports whether faults hold the given condition for the site.
|
|
func hasFault(faults []Fault, condition, site string) bool {
|
|
for _, f := range faults {
|
|
if f.Condition == condition && f.Site == site {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// countEligible routes the eligible count through the test seam when one is
|
|
// set, else the store.
|
|
func (p *Poller) countEligible(site string) (int, error) {
|
|
if p.eligibleCount != nil {
|
|
return p.eligibleCount(site)
|
|
}
|
|
return p.Store.EligibleSeriesCount(site)
|
|
}
|
|
|
|
// 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
|
|
}
|
|
|
|
// anyForced reports whether the due list holds a forced Series: one whose
|
|
// owner check-now request (issue #146) has not been answered yet. A human
|
|
// asking wakes a sleeping Chrome even below the wake thresholds; the request
|
|
// itself still ages visibly if the home machine is off.
|
|
func anyForced(due []store.Series) bool {
|
|
for _, sr := range due {
|
|
if sr.Forced {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// 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
|
|
}
|
|
// recordFailure keeps one Series' failure row in step with its read
|
|
// (ADR-0016): a failure word is upserted, a successful read deletes the row,
|
|
// and refused or unreachable issue no statement at all. Called for every
|
|
// outcome from the pass loop, so the four failure words and the success path
|
|
// share one write point, and a forced Poll that reads the page clears through
|
|
// the ordinary success path — no branch of its own. Log a store failure and
|
|
// carry on: this is best-effort, and no single bad Series may stall a Lane.
|
|
func (p *Poller) recordFailure(sr store.Series, outcome readOutcome) {
|
|
if outcome == outcomeSuccess {
|
|
if err := p.Store.ClearSeriesFailure(sr.Site, sr.SeriesID); err != nil {
|
|
log.Printf("latest poll %q: clear failure: %v", sr.Key(), err)
|
|
}
|
|
return
|
|
}
|
|
word := outcome.word()
|
|
if word == "" {
|
|
return
|
|
}
|
|
if err := p.Store.RecordSeriesFailure(sr.Site, sr.SeriesID, word, p.Now().UnixMilli()); err != nil {
|
|
log.Printf("latest poll %q: record failure: %v", sr.Key(), err)
|
|
}
|
|
}
|
|
|
|
// 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 outcome classifies the read for
|
|
// the pass row (issue #141), so the Lane can count a refusal, a lost browser,
|
|
// a chapter-less page, an unfetchable address, a missing page or a transport
|
|
// error without re-deriving the taxonomy.
|
|
func (p *Poller) checkOne(ctx context.Context, sr store.Series) (outcome readOutcome) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
log.Printf("latest poll %q: recovered from panic: %v", sr.Key(), r)
|
|
outcome = outcomeError
|
|
}
|
|
}()
|
|
|
|
// 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 outcomeError
|
|
}
|
|
|
|
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 outcomeUnfetchable
|
|
case errors.Is(err, errNoFetcher):
|
|
log.Printf("latest poll %q: no fetcher for site %q", sr.Key(), sr.Site)
|
|
return outcomeUnfetchable
|
|
}
|
|
// 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)
|
|
if errors.Is(err, errChallengeHeld) {
|
|
return outcomeRefused
|
|
}
|
|
if errors.Is(err, errBrowserInterrupted) {
|
|
return outcomeUnreachable
|
|
}
|
|
if errors.Is(err, errNotFound) {
|
|
return outcomeNotFound
|
|
}
|
|
return outcomeError
|
|
}
|
|
// 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.
|
|
// A forced pass writes the Cover through the replace path instead.
|
|
if sr.Forced {
|
|
p.replaceCover(ctx, sr, facts.Cover)
|
|
} else {
|
|
p.fillBlankCover(ctx, sr, facts.Cover)
|
|
}
|
|
|
|
// Learned from the successful read: the write sits after the error switch
|
|
// (a refused, unreachable or errored read reaches nothing) and before the
|
|
// returns below — a completed page whose chapter number did not change
|
|
// still has to write. The transition is zero-versus-nonzero, not the
|
|
// stamp's value: a Series still completed keeps its original stamp, so the
|
|
// age #170 prints is "since the Site first said so"; one that stopped
|
|
// being completed is zeroed.
|
|
stamp := int64(0)
|
|
if facts.SiteCompleted {
|
|
stamp = p.Now().UnixMilli()
|
|
}
|
|
if (sr.SiteCompletedAt == 0) != (stamp == 0) {
|
|
if err := p.Store.SetSiteCompletedAt(sr.Site, sr.SeriesID, stamp); err != nil {
|
|
// Best-effort, like every poller write: never change the outcome
|
|
// word the pass counts.
|
|
log.Printf("latest poll %q: set site completed: %v", sr.Key(), err)
|
|
}
|
|
}
|
|
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 outcomeNoChapter
|
|
}
|
|
|
|
// 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 outcomeSuccess
|
|
}
|
|
|
|
// 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 outcomeError
|
|
}
|
|
log.Printf("latest poll %q: latest is now %s", sr.Key(), facts.Latest.Label)
|
|
return outcomeSuccess
|
|
}
|
|
|
|
// 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 {
|
|
// The report is neither confirmed nor contradicted, but it is answered:
|
|
// the value about to be stored is the Poll's own, so leaving the
|
|
// attribution would credit this Reader with the next Poll's agreement
|
|
// and blame them if the Site later retracts.
|
|
if err := p.Store.ClearSightingAttribution(sr.Site, sr.SeriesID, *sr.LatestRaisedBy); err != nil {
|
|
log.Printf("latest poll %q: clear sighting attribution: %v", sr.Key(), err)
|
|
}
|
|
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.
|
|
//
|
|
// The owner's series URL repair (issue #151) is a second caller: the web
|
|
// layer validates with this same gate before storing a repair, so there is
|
|
// never a second copy of it.
|
|
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
|
|
}
|