feat(latest): record one poll pass per exit with skip reason and outcome counts
Every way a Lane pass can end now writes exactly one durable row: a skip value naming the exit (paused, refusing, sidecar-down, no-fetcher, due-query, asleep, eligible-count, nothing-eligible, or empty for the loop), five outcome counts from the classification the Series read already makes, and carry-forward of the previous pass's figures exactly when the pass's own gap is zero. Refusal is durable through the poll_lanes row, so a restart does not re-probe a Site inside its backoff. Retention is 14 days. The mid-loop browser-unreachable return writes an empty skip by design: a tenth value is not invented here. (#141)
This commit is contained in:
@@ -48,6 +48,12 @@ type Poller struct {
|
||||
// same failure-isolated prefetch path.
|
||||
CoverBytesFetch CoverBytesFetcher
|
||||
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).
|
||||
@@ -216,6 +222,74 @@ func (p *Poller) runOnce(ctx context.Context) {
|
||||
}
|
||||
}
|
||||
|
||||
// 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 (
|
||||
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, and errors everything else.
|
||||
type readOutcome int
|
||||
|
||||
const (
|
||||
outcomeSuccess readOutcome = iota
|
||||
outcomeRefused
|
||||
outcomeUnreachable
|
||||
outcomeNoChapter
|
||||
outcomeUnfetchable
|
||||
outcomeError
|
||||
)
|
||||
|
||||
// outcomeCounts are the five named outcome counts of one pass. A success
|
||||
// count is derived, never stored: checked minus the four, with unreachable
|
||||
// excluded because the sidecar-loss path returns before the checked counter
|
||||
// increments (issue #141).
|
||||
type outcomeCounts struct {
|
||||
refused, unreachable, noChapter, unfetchable, errors 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 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
|
||||
@@ -226,9 +300,32 @@ func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.
|
||||
// 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) {
|
||||
// Durable pass log (issue #141): one row per exit, written from the same
|
||||
// snapshot so the two recordings cannot disagree.
|
||||
rec := passRecord{site: name, ranAt: now.UnixMilli()}
|
||||
defer func() { p.recordPass(rec, st) }()
|
||||
|
||||
// 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.
|
||||
return until.Sub(now)
|
||||
rec.skip = skipRefusing
|
||||
return time.Duration(refuseUntil-now.UnixMilli()) * time.Millisecond
|
||||
}
|
||||
if isBrowserSite(name) {
|
||||
if downFor, down := p.browserDownFor(now); down && downFor < refuseBackoff {
|
||||
@@ -236,6 +333,7 @@ func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.
|
||||
// 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
|
||||
}
|
||||
@@ -246,6 +344,7 @@ func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.
|
||||
// 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
|
||||
st.Gap = defaultGap
|
||||
return defaultGap
|
||||
}
|
||||
@@ -253,6 +352,7 @@ func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.
|
||||
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)
|
||||
st.Gap = defaultGap
|
||||
return defaultGap
|
||||
@@ -263,6 +363,7 @@ func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.
|
||||
// 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.
|
||||
rec.skip = skipAsleep
|
||||
st.Gap, st.Asleep = defaultGap, true
|
||||
return defaultGap
|
||||
}
|
||||
@@ -276,8 +377,9 @@ func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.
|
||||
}
|
||||
}
|
||||
|
||||
eligible, err := p.Store.EligibleSeriesCount(name)
|
||||
eligible, err := p.countEligible(name)
|
||||
if err != nil {
|
||||
rec.skip = skipEligibleCount
|
||||
log.Printf("latest poll %s: eligible count: %v", name, err)
|
||||
st.Gap = defaultGap
|
||||
return defaultGap
|
||||
@@ -290,6 +392,7 @@ func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.
|
||||
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
|
||||
}
|
||||
|
||||
@@ -314,20 +417,24 @@ func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.
|
||||
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
|
||||
}
|
||||
outcome := p.checkOne(ctx, sr)
|
||||
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)
|
||||
st.Checked++
|
||||
}
|
||||
if st.Checked > 0 {
|
||||
@@ -335,16 +442,65 @@ func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.
|
||||
}
|
||||
if refusals >= 2 {
|
||||
p.setRefusalBackoff(name, now.Add(refuseBackoff))
|
||||
// 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. The
|
||||
// in-memory twin is still written because the Lane status block reads
|
||||
// it directly; the gate reads the durable stamp, so a restart does not
|
||||
// forget the refusal.
|
||||
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
|
||||
}
|
||||
|
||||
func (p *Poller) refusalBackoff(name string) time.Time {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
return p.refuseUntil[name]
|
||||
// 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.
|
||||
func (p *Poller) recordPass(rec passRecord, st LaneState) {
|
||||
row := store.LanePass{
|
||||
Site: rec.site,
|
||||
RanAt: rec.ranAt,
|
||||
Skip: rec.skip,
|
||||
Due: st.Due,
|
||||
Checked: st.Checked,
|
||||
GapMS: st.Gap.Milliseconds(),
|
||||
Clamped: st.Clamped,
|
||||
Refused: rec.counts.refused,
|
||||
Unreachable: rec.counts.unreachable,
|
||||
NoChapter: rec.counts.noChapter,
|
||||
Unfetchable: rec.counts.unfetchable,
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
|
||||
func (p *Poller) setRefusalBackoff(name string, until time.Time) {
|
||||
@@ -409,13 +565,15 @@ func maxSeriesWait(due []store.Series, now time.Time, rest time.Duration) time.D
|
||||
|
||||
// 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 {
|
||||
// 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 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
|
||||
}
|
||||
}()
|
||||
|
||||
@@ -427,7 +585,7 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) error {
|
||||
// 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
|
||||
return outcomeError
|
||||
}
|
||||
|
||||
facts, err := readSeriesPage(ctx, sr.Site, sr.SeriesURL, p.BrowserFetch, p.Fetch)
|
||||
@@ -438,17 +596,23 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) error {
|
||||
// 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
|
||||
return outcomeUnfetchable
|
||||
case errors.Is(err, errNoFetcher):
|
||||
log.Printf("latest poll %q: no fetcher for site %q", sr.Key(), sr.Site)
|
||||
return err
|
||||
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)
|
||||
return err
|
||||
if errors.Is(err, errChallengeHeld) {
|
||||
return outcomeRefused
|
||||
}
|
||||
if errors.Is(err, errBrowserInterrupted) {
|
||||
return outcomeUnreachable
|
||||
}
|
||||
return outcomeError
|
||||
}
|
||||
// A legacy cover source is healed independently of the page read.
|
||||
p.healCover(ctx, sr)
|
||||
@@ -459,7 +623,7 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) error {
|
||||
// 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
|
||||
return outcomeNoChapter
|
||||
}
|
||||
|
||||
// The Poll is the oracle for whatever Sighting last raised this Series
|
||||
@@ -472,7 +636,7 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) error {
|
||||
// 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
|
||||
return outcomeSuccess
|
||||
}
|
||||
|
||||
// Series-level write: the row is shared, so one update refreshes every
|
||||
@@ -481,10 +645,10 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) error {
|
||||
// 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
|
||||
return outcomeError
|
||||
}
|
||||
log.Printf("latest poll %q: latest is now %s", sr.Key(), facts.Latest.Label)
|
||||
return nil
|
||||
return outcomeSuccess
|
||||
}
|
||||
|
||||
// judgeSighting settles the Sighting the Series' stored Latest Chapter is owed
|
||||
|
||||
Reference in New Issue
Block a user