Spec #134, all ten tickets. Closes #134. ## What ships The admin surface becomes four bookmarkable addresses behind one nav row, and Lane observability stops dying with the process. - **#138** `/admin` splits into Overview, Lanes, Readers, Series, each a real route with the active tab underlined. - **#139** `poll_passes` and `poll_lanes` land as durable tables with their store surface. - **#140** cross-Series admin read model, with the privacy boundary in the projection: the Reader id that raised a Latest Chapter never leaves the store package. - **#141** the poller records exactly one pass row per exit, with a skip reason and outcome counts. - **#142** Series list: eight hygiene filters, Site and Library narrowing, paging — all of it in the query string, so a filtered list is a bookmark. - **#143** Overview: a three-state verdict line and a stats block where every non-zero figure links to the list that counts it. - **#144** per-Series detail page, keyed by the `site:series_id` composite the rest of the system already uses. - **#145** the Lanes page reads the database; the in-memory Lane state, `web.LaneReporter` and `latest.Status` are deleted. - **#146** Forced Poll: *Check now* stamps `series.force_poll_at` and never commands the poller. - **#147** pause and resume one Site's Lane, with a mandatory 1h/6h/24h expiry. ## Shape of the design Two decisions carry the rest. **Commands go through the database, never at the poller**: both *Check now* and a Lane pause write a row the next pass reads, so they survive a restart and the whole surface stays testable with no poller running. And **pending is derived, never stored** — the request stamp being newer than the check stamp — which self-clears on the check stamp with no second write and no sweeper, because the check stamp is written before the fetch. ADRs: `docs/adr/0012-persisted-lane-state.md`, `docs/adr/0013-commands-through-the-database.md`. ## Verification `go test ./...` green on the merged base (`264839e`), all packages, Docker-backed. `gofmt -l` and `go vet` clean. Every ticket was reviewed on both axes (`cr-spec` + `cr-standards`) before merge. ## Known, non-blocking - **#143** the verdict ignores never-reported Lanes when other Lanes have reported, and the per-Site table lists Sites that have Series rather than the whole registry. The ticket prose asks for eight hygiene figures per Site; the design mock and the landed `.tbl.sites` grid both say six columns, and the mock won. - **#146** two `SeriesPage` scans per press instead of a keyed read — `ponytail:`-commented in-tree with the upgrade path. - **#147** a paused Site with no pass row yet renders no row and so no control, since the Lanes page lists Sites that have passed. - **#141** a sibling browser Lane declining at the top of a pass records as `sidecar-down`. Specified deliberately; the later spec in this series settles it. Reviewed-on: #148 Co-authored-by: Sulthan Zaki <sultankiki05@gmail.com> Co-committed-by: Sulthan Zaki <sultankiki05@gmail.com>
This commit was merged in pull request #148.
This commit is contained in:
@@ -135,6 +135,33 @@ func TestAcquireFillsChapterAndCoverFromOneFetch(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// A pause governs the Lane only: a Reader's first bookmark of a Series on a
|
||||
// paused Site still reads the page, because acquisition is the creation-time
|
||||
// fetch, not the poll queue (issue #147).
|
||||
func TestAcquireIgnoresLanePause(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
if err := s.PauseLane("asura", time.Now().Add(6*time.Hour).UnixMilli()); err != nil {
|
||||
t.Fatalf("PauseLane: %v", err)
|
||||
}
|
||||
page := &fakeFetcher{body: asuraSeriesAndCoverFixture, status: 200}
|
||||
covers := &fakeBytesCoverFetcher{body: []byte("cover-bytes"), contentType: "image/jpeg"}
|
||||
acq := newAcquirer(s, page, covers)
|
||||
|
||||
bookmarkNewSeries(t, s, acquireSeriesURL)
|
||||
acq.Wait()
|
||||
|
||||
if got := page.callCount(); got != 1 {
|
||||
t.Fatalf("series page fetches on a paused Site = %d, want 1", got)
|
||||
}
|
||||
if got := covers.callCount(); got != 1 {
|
||||
t.Fatalf("cover fetches = %d, want 1", got)
|
||||
}
|
||||
got := readBookmark(t, s, acquireKey)
|
||||
if got.LatestChapterNum == nil || *got.LatestChapterNum != 181 {
|
||||
t.Fatalf("LatestChapterNum = %v, want 181", got.LatestChapterNum)
|
||||
}
|
||||
}
|
||||
|
||||
// A Series that already exists is not re-acquired: no fetch, and the Cover it
|
||||
// already has is left alone.
|
||||
func TestAcquireSkipsAnExistingSeries(t *testing.T) {
|
||||
|
||||
@@ -5,7 +5,6 @@ import (
|
||||
"errors"
|
||||
"log"
|
||||
"net/url"
|
||||
"sort"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
@@ -48,19 +47,22 @@ 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).
|
||||
// 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
|
||||
// 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.
|
||||
@@ -167,14 +169,7 @@ func fetcherFor(site string, browser, tls Fetcher) Fetcher {
|
||||
// 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 laneNames() []string { return SiteNames() }
|
||||
|
||||
func (p *Poller) Run(ctx context.Context) {
|
||||
names := laneNames()
|
||||
@@ -216,28 +211,120 @@ 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 (
|
||||
// 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, 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
|
||||
// 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) {
|
||||
// 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(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.
|
||||
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 {
|
||||
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
|
||||
// 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
|
||||
return RefuseBackoff - downFor
|
||||
}
|
||||
}
|
||||
s := sites[name]
|
||||
@@ -246,24 +333,29 @@ 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).
|
||||
st.Gap = defaultGap
|
||||
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)
|
||||
st.Gap = defaultGap
|
||||
fig.Gap = defaultGap
|
||||
return defaultGap
|
||||
}
|
||||
st.Due = len(due)
|
||||
if s.Browser != nil && f == p.BrowserFetch && !browserWakeDue(due, now, s.Rest) {
|
||||
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. The Lane still paces at the default gap, which is what
|
||||
// the owner's page must show rather than a zero.
|
||||
st.Gap, st.Asleep = defaultGap, true
|
||||
// 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 {
|
||||
@@ -276,20 +368,22 @@ 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
|
||||
fig.Gap = defaultGap
|
||||
return defaultGap
|
||||
}
|
||||
gap, clamped := effectiveGap(s, eligible)
|
||||
st.Gap, st.Clamped = gap, clamped
|
||||
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
|
||||
}
|
||||
|
||||
@@ -300,7 +394,7 @@ func (p *Poller) runLanePass(ctx context.Context, name string, paced bool) time.
|
||||
}
|
||||
if refusals >= 2 {
|
||||
// This Site refused twice in a row: the remaining Series are left
|
||||
// unstamped and due, and the Lane waits refuseBackoff before
|
||||
// unstamped and due, and the Lane waits RefuseBackoff before
|
||||
// trying it again.
|
||||
break
|
||||
}
|
||||
@@ -314,46 +408,98 @@ 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
|
||||
}
|
||||
st.Checked++
|
||||
rec.counts.add(outcome)
|
||||
fig.Checked++
|
||||
}
|
||||
if st.Checked > 0 {
|
||||
log.Printf("latest poll %s: due=%d checked=%d", name, len(due), st.Checked)
|
||||
if fig.Checked > 0 {
|
||||
log.Printf("latest poll %s: due=%d checked=%d", name, len(due), fig.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
|
||||
// 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
|
||||
}
|
||||
|
||||
func (p *Poller) refusalBackoff(name string) time.Time {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
return p.refuseUntil[name]
|
||||
// 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
|
||||
}
|
||||
|
||||
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)
|
||||
// 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, 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,
|
||||
Errors: rec.counts.errors,
|
||||
}
|
||||
p.refuseUntil[name] = until
|
||||
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)
|
||||
}
|
||||
|
||||
// setBrowserDown records when a browser Lane lost the sidecar. It is Poller
|
||||
@@ -366,7 +512,7 @@ func (p *Poller) setBrowserDown(now time.Time) {
|
||||
|
||||
// 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
|
||||
// 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()
|
||||
@@ -395,6 +541,19 @@ func browserWakeDue(due []store.Series, now time.Time, rest time.Duration) bool
|
||||
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 {
|
||||
@@ -409,13 +568,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 +588,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 +599,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 +626,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 +639,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 +648,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
|
||||
|
||||
@@ -427,9 +427,12 @@ func TestRunLogsLaneDefaults(t *testing.T) {
|
||||
log.SetOutput(&logs)
|
||||
t.Cleanup(func() { log.SetOutput(previous) })
|
||||
|
||||
// The pass log writes through the store, so Run needs a real one — a
|
||||
// poller without a Store is not a poller.
|
||||
s, _ := newTestStore(t)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
(&Poller{Now: func() time.Time { return time.Now() }}).Run(ctx)
|
||||
(&Poller{Store: s, Now: func() time.Time { return time.Now() }}).Run(ctx)
|
||||
|
||||
got := logs.String()
|
||||
for _, want := range []string{"6 lanes", "rest=1h0m0s", "gap=10s"} {
|
||||
@@ -1305,7 +1308,7 @@ func TestRunOnceOrdersBySharednessThenAge(t *testing.T) {
|
||||
}
|
||||
|
||||
// Two refusals in one pass stop the Lane: the remaining Series stay unstamped
|
||||
// and due, and the Lane backs off for refuseBackoff before trying the Site
|
||||
// and due, and the Lane backs off for RefuseBackoff before trying the Site
|
||||
// again. One hostile Site burns only its own Lane's budget (issue #100).
|
||||
func TestRunOnceSiteRefusalSkipsRestOfLaneAndBacksOff(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
@@ -1388,11 +1391,9 @@ func TestBrowserLaneWakeThresholds(t *testing.T) {
|
||||
if got := browser.callCount(); got != 0 {
|
||||
t.Fatalf("browser fetches with 3 freshly-due series = %d, want 0 (Chrome stays asleep)", got)
|
||||
}
|
||||
// The owner's page reads this state off the snapshot, and Due-without-
|
||||
// Checked has to be distinguishable there from a Lane that has stopped.
|
||||
if lane := laneByName(t, p, "kagane"); !lane.Asleep || lane.Due != 3 || lane.Checked != 0 {
|
||||
t.Fatalf("asleep kagane lane = %+v, want Asleep with 3 due and 0 checked", lane)
|
||||
}
|
||||
// The skipped pass still records its row, carrying the due count it never
|
||||
// read; the Lanes page (issue #145) reads that row — see
|
||||
// TestRunLanePassRecordsEveryExit/"asleep".
|
||||
// 5 due crosses the count threshold.
|
||||
for i := 3; i < 5; i++ {
|
||||
seed(i)
|
||||
@@ -1401,9 +1402,6 @@ func TestBrowserLaneWakeThresholds(t *testing.T) {
|
||||
if got := browser.callCount(); got != 5 {
|
||||
t.Fatalf("browser fetches with 5 due series = %d, want 5", got)
|
||||
}
|
||||
if lane := laneByName(t, p, "kagane"); lane.Asleep {
|
||||
t.Fatalf("woken kagane lane still reports Asleep: %+v", lane)
|
||||
}
|
||||
// A single long-neglected series wakes the browser by age alone.
|
||||
seedForCheck(t, s, "kagane:ancient", "https://kagane.to/series/ancient", 0)
|
||||
p.runOnce(context.Background())
|
||||
@@ -1412,19 +1410,6 @@ func TestBrowserLaneWakeThresholds(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// laneByName pulls one Lane out of the poller's snapshot, failing rather than
|
||||
// returning a zero LaneState a caller would assert against by accident.
|
||||
func laneByName(t *testing.T, p *Poller, site string) LaneState {
|
||||
t.Helper()
|
||||
for _, lane := range p.LaneStatus().Lanes {
|
||||
if lane.Site == site {
|
||||
return lane
|
||||
}
|
||||
}
|
||||
t.Fatalf("no %q lane in the snapshot", site)
|
||||
return LaneState{}
|
||||
}
|
||||
|
||||
// When one browser Lane loses the sidecar, the round's remaining browser
|
||||
// Lanes are skipped: every fetch would fail anyway, and their Series must not
|
||||
// burn their stamps on a dead Chrome (issue #100).
|
||||
@@ -1456,7 +1441,7 @@ func TestRunOnceUnreachableBrowserStopsBrowserLanes(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// The shared flag decays after refuseBackoff: the next round probes
|
||||
// The shared flag decays after RefuseBackoff: the next round probes
|
||||
// again. comix's Series is resting (stamped last round), so the probe
|
||||
// falls to kagane — the only Lane with something due — and its fresh
|
||||
// loss re-gates the Lanes behind it.
|
||||
@@ -1602,19 +1587,21 @@ func TestRunOnceClampWarningNamesTheSite(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// The owner's admin page (issue #102) reads Lane state out of the poller.
|
||||
// Before any pass the snapshot is empty — a restart must render "no data
|
||||
// yet", not zeroes — and each pass records what it saw: the frozen clock,
|
||||
// the due count and the pace, with refusal backoff derived at snapshot time.
|
||||
func TestLaneStatus(t *testing.T) {
|
||||
// The Lanes page (issue #145) reads the durable pass log, so the poller's
|
||||
// only duty to it is that every return path writes its row; the snapshot it
|
||||
// used to mirror into memory is gone. A refusing pass still records why it
|
||||
// declined, and the latest row per Site is what the page renders — nothing
|
||||
// else is left to assert against here, because the page's seam moved into the
|
||||
// web layer (web_test.go, seeded-row tests).
|
||||
func TestPassLogIsTheLanesPagesOnlyWindow(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
now := time.UnixMilli(5_000_000)
|
||||
p := newTestPoller(t, s, &fakeFetcher{body: asuraSeriesFixture, status: 200}, now)
|
||||
|
||||
if st := p.LaneStatus(); len(st.Lanes) != 0 {
|
||||
t.Fatalf("lanes before any pass = %d, want 0 (nothing has run)", len(st.Lanes))
|
||||
} else if st.BrowserConfigured || st.BrowserReachable {
|
||||
t.Fatalf("browser before any pass = configured=%v reachable=%v, want false without a browser fetcher", st.BrowserConfigured, st.BrowserReachable)
|
||||
// No pass has run: the log is empty, which the page renders as "none
|
||||
// observed" rather than confident zeroes.
|
||||
if _, ok, err := s.LatestLanePass("asura"); err != nil || ok {
|
||||
t.Fatalf("LatestLanePass before any pass = ok %v err %v, want no row", ok, err)
|
||||
}
|
||||
|
||||
seedForCheck(t, s, "asura:chronicles", "https://asurascans.com/series/chronicles", 0)
|
||||
@@ -1627,65 +1614,726 @@ func TestLaneStatus(t *testing.T) {
|
||||
}
|
||||
p.runOnce(context.Background())
|
||||
|
||||
st := p.LaneStatus()
|
||||
var asura, kagane LaneState
|
||||
for _, lane := range st.Lanes {
|
||||
switch lane.Site {
|
||||
case "asura":
|
||||
asura = lane
|
||||
case "kagane":
|
||||
kagane = lane
|
||||
}
|
||||
// Both lanes wrote their rows: asura the pass it read; kagane the
|
||||
// mid-loop double-refusal exit, which records an empty skip by design
|
||||
// with its two refused counts — the durable row is the whole record, and
|
||||
// the page reads it rather than a snapshot.
|
||||
asura := latestPassFor(t, s, "asura")
|
||||
if asura.Due != 1 || asura.Checked != 1 || asura.GapMS == 0 {
|
||||
t.Fatalf("asura row = %+v, want due 1 checked 1 with the Lane's pace", asura)
|
||||
}
|
||||
if st.Lanes[0].Site != "asura" {
|
||||
t.Fatalf("first lane = %q, want asura (snapshot sorted by Site)", st.Lanes[0].Site)
|
||||
}
|
||||
if asura.Site == "" {
|
||||
t.Fatalf("asura missing from snapshot: %+v", st.Lanes)
|
||||
}
|
||||
if !asura.LastRun.Equal(now) {
|
||||
t.Fatalf("asura LastRun = %s, want the frozen clock %s", asura.LastRun, now)
|
||||
}
|
||||
if asura.Due != 1 {
|
||||
t.Fatalf("asura Due = %d, want 1", asura.Due)
|
||||
}
|
||||
if asura.Gap == 0 {
|
||||
t.Fatal("asura Gap = 0, want the Lane's pace")
|
||||
}
|
||||
if asura.Browser || asura.Refusing {
|
||||
t.Fatalf("asura = %+v, want a TLS Lane that is not refusing", asura)
|
||||
}
|
||||
if !kagane.Browser || !kagane.Refusing {
|
||||
t.Fatalf("kagane = %+v, want a browser Lane in refusal backoff", kagane)
|
||||
}
|
||||
if !st.BrowserConfigured || !st.BrowserReachable {
|
||||
t.Fatalf("browser after round = configured=%v reachable=%v, want true/true (sidecar never lost)", st.BrowserConfigured, st.BrowserReachable)
|
||||
}
|
||||
if asura.Checked != 1 {
|
||||
t.Fatalf("asura Checked = %d, want the one Series it read", asura.Checked)
|
||||
kagane := latestPassFor(t, s, "kagane")
|
||||
if kagane.Skip != "" || kagane.Refused != 2 || kagane.Due != 2 || kagane.GapMS == 0 {
|
||||
t.Fatalf("kagane row = %+v, want an empty-skip double-refusal pass with 2 refused", kagane)
|
||||
}
|
||||
|
||||
// A pass that declines to look (kagane is now in backoff) must not restate
|
||||
// the figures it never gathered as zeroes: the last real pass's due count
|
||||
// and pace stand until a pass replaces them.
|
||||
// With the refusal now durable, the next pass declines ahead of the loop:
|
||||
// it records the refusing skip and carries the previous pass's figures
|
||||
// forward rather than restating zeroes it never gathered (the carry logic
|
||||
// itself is driven in TestRunLanePassCarryForwardOnlyWhenGapZero).
|
||||
before := kagane
|
||||
if before.Due == 0 || before.Gap == 0 {
|
||||
t.Fatalf("kagane after its refusing pass = %+v, want the figures that pass gathered", before)
|
||||
}
|
||||
p.Now = func() time.Time { return now.Add(time.Minute) }
|
||||
p.runOnce(context.Background())
|
||||
for _, lane := range p.LaneStatus().Lanes {
|
||||
if lane.Site != "kagane" {
|
||||
continue
|
||||
}
|
||||
if lane.Due != before.Due || lane.Gap != before.Gap {
|
||||
t.Fatalf("kagane after a skipped pass = due %d gap %s, want the previous pass's %d / %s",
|
||||
lane.Due, lane.Gap, before.Due, before.Gap)
|
||||
}
|
||||
again := latestPassFor(t, s, "kagane")
|
||||
if again.Skip != SkipRefusing {
|
||||
t.Fatalf("kagane refusing pass skip = %q, want %q", again.Skip, SkipRefusing)
|
||||
}
|
||||
|
||||
// A lost sidecar reads as unreachable for the same window the Lanes skip.
|
||||
p.setBrowserDown(now)
|
||||
if st := p.LaneStatus(); !st.BrowserConfigured || st.BrowserReachable {
|
||||
t.Fatalf("browser after loss = configured=%v reachable=%v, want true/false", st.BrowserConfigured, st.BrowserReachable)
|
||||
if again.Due != before.Due || again.GapMS != before.GapMS || again.Checked != before.Checked {
|
||||
t.Fatalf("kagane after a skipped pass = due %d gap %d, want the previous pass's %d / %d",
|
||||
again.Due, again.GapMS, before.Due, before.GapMS)
|
||||
}
|
||||
}
|
||||
|
||||
// latestPassFor reads a Site's newest durable pass row, failing rather than
|
||||
// returning a zero LanePass a caller would assert against by accident.
|
||||
func latestPassFor(t *testing.T, s *store.Store, site string) store.LanePass {
|
||||
t.Helper()
|
||||
pass, ok, err := s.LatestLanePass(site)
|
||||
if err != nil || !ok {
|
||||
t.Fatalf("LatestLanePass(%s): ok=%v err=%v", site, ok, err)
|
||||
}
|
||||
return pass
|
||||
}
|
||||
|
||||
// countPassRows counts a Site's durable pass rows, for asserting that a pass
|
||||
// records exactly one.
|
||||
func countPassRows(t *testing.T, dbURL, site string) int {
|
||||
t.Helper()
|
||||
db, err := sql.Open("pgx", dbURL)
|
||||
if err != nil {
|
||||
t.Fatalf("open %s: %v", dbURL, err)
|
||||
}
|
||||
defer db.Close()
|
||||
var n int
|
||||
if err := db.QueryRow(`SELECT count(*) FROM poll_passes WHERE site = $1`, site).Scan(&n); err != nil {
|
||||
t.Fatalf("count passes for %s: %v", site, err)
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
// One durable pass row per exit, with the skip value naming the exit. The
|
||||
// mid-loop browser-unreachable return also writes one row — but with an empty
|
||||
// skip, so the row is stall-shaped (due > 0, checked 0, skip ”), matching the
|
||||
// deliberate absence of a tenth skip value.
|
||||
func TestRunLanePassRecordsEveryExit(t *testing.T) {
|
||||
now := time.UnixMilli(5_000_000)
|
||||
|
||||
t.Run("paused", func(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
seedForCheck(t, s, "asura:x", "https://asurascans.com/series/x", 0)
|
||||
if err := s.PauseLane("asura", now.Add(30*time.Minute).UnixMilli()); err != nil {
|
||||
t.Fatalf("PauseLane: %v", err)
|
||||
}
|
||||
f := &fakeFetcher{body: asuraSeriesFixture, status: 200}
|
||||
p := newTestPoller(t, s, f, now)
|
||||
if pace := p.runLanePass(context.Background(), "asura", false); pace != 30*time.Minute {
|
||||
t.Fatalf("paused pace = %s, want 30m (sleep until the expiry)", pace)
|
||||
}
|
||||
if f.callCount() != 0 {
|
||||
t.Fatalf("fetches while paused = %d, want 0", f.callCount())
|
||||
}
|
||||
if got := readLatestCheckedAt(t, s, "asura:x"); got != 0 {
|
||||
t.Fatalf("stamp while paused = %d, want 0 (Series stay due and unstamped)", got)
|
||||
}
|
||||
if pass := latestPassFor(t, s, "asura"); pass.Skip != SkipPaused {
|
||||
t.Fatalf("skip = %q, want %q", pass.Skip, SkipPaused)
|
||||
}
|
||||
if got := countPassRows(t, dbURL, "asura"); got != 1 {
|
||||
t.Fatalf("pass rows = %d, want exactly 1", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("refusing", func(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
seedForCheck(t, s, "kagane:x", "https://kagane.to/series/x", 0)
|
||||
if err := s.SetLaneRefusal("kagane", now.Add(10*time.Minute).UnixMilli()); err != nil {
|
||||
t.Fatalf("SetLaneRefusal: %v", err)
|
||||
}
|
||||
browser := &fakeFetcher{status: 200}
|
||||
p := newTestPoller(t, s, &fakeFetcher{status: 200}, now)
|
||||
p.BrowserFetch = browser
|
||||
if pace := p.runLanePass(context.Background(), "kagane", false); pace != 10*time.Minute {
|
||||
t.Fatalf("refusing pace = %s, want 10m", pace)
|
||||
}
|
||||
if browser.callCount() != 0 {
|
||||
t.Fatalf("fetches while refusing = %d, want 0", browser.callCount())
|
||||
}
|
||||
if pass := latestPassFor(t, s, "kagane"); pass.Skip != SkipRefusing {
|
||||
t.Fatalf("skip = %q, want %q", pass.Skip, SkipRefusing)
|
||||
}
|
||||
if got := countPassRows(t, dbURL, "kagane"); got != 1 {
|
||||
t.Fatalf("pass rows = %d, want exactly 1", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("sidecar-down", func(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
seedForCheck(t, s, "kagane:x", "https://kagane.to/series/x", 0)
|
||||
browser := &fakeFetcher{body: kaganeAPIFixture, status: 200}
|
||||
p := newTestPoller(t, s, &fakeFetcher{status: 200}, now)
|
||||
p.BrowserFetch = browser
|
||||
p.setBrowserDown(now) // a sibling Lane lost Chrome within the backoff window
|
||||
if pace := p.runLanePass(context.Background(), "kagane", false); pace != RefuseBackoff {
|
||||
t.Fatalf("sidecar-down pace = %s, want %s", pace, RefuseBackoff)
|
||||
}
|
||||
if browser.callCount() != 0 {
|
||||
t.Fatalf("fetches with the sidecar down = %d, want 0", browser.callCount())
|
||||
}
|
||||
if pass := latestPassFor(t, s, "kagane"); pass.Skip != SkipSidecarDown {
|
||||
t.Fatalf("skip = %q, want %q", pass.Skip, SkipSidecarDown)
|
||||
}
|
||||
if got := countPassRows(t, dbURL, "kagane"); got != 1 {
|
||||
t.Fatalf("pass rows = %d, want exactly 1", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("no-fetcher", func(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
seedForCheck(t, s, "comix:c", "https://comix.to/title/c", 0)
|
||||
p := newTestPoller(t, s, &fakeFetcher{status: 200}, now)
|
||||
if pace := p.runLanePass(context.Background(), "comix", false); pace != defaultGap {
|
||||
t.Fatalf("no-fetcher pace = %s, want %s", pace, defaultGap)
|
||||
}
|
||||
pass := latestPassFor(t, s, "comix")
|
||||
if pass.Skip != SkipNoFetcher || pass.GapMS != defaultGap.Milliseconds() {
|
||||
t.Fatalf("no-fetcher pass = %+v, want skip %q with its own gap %s", pass, SkipNoFetcher, defaultGap)
|
||||
}
|
||||
if got := countPassRows(t, dbURL, "comix"); got != 1 {
|
||||
t.Fatalf("pass rows = %d, want exactly 1", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("due-query", func(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
seedForCheck(t, s, "asura:x", "https://asurascans.com/series/x", 0)
|
||||
// The due query is the pass's first store read after the gates; making
|
||||
// it fail without touching poll_passes takes its tables away.
|
||||
db, err := sql.Open("pgx", dbURL)
|
||||
if err != nil {
|
||||
t.Fatalf("open %s: %v", dbURL, err)
|
||||
}
|
||||
if _, err := db.Exec(`DROP TABLE bookmarks`); err != nil {
|
||||
t.Fatalf("drop bookmarks: %v", err)
|
||||
}
|
||||
db.Close()
|
||||
p := newTestPoller(t, s, &fakeFetcher{body: asuraSeriesFixture, status: 200}, now)
|
||||
if pace := p.runLanePass(context.Background(), "asura", false); pace != defaultGap {
|
||||
t.Fatalf("due-query pace = %s, want %s", pace, defaultGap)
|
||||
}
|
||||
if pass := latestPassFor(t, s, "asura"); pass.Skip != SkipDueQuery {
|
||||
t.Fatalf("skip = %q, want %q", pass.Skip, SkipDueQuery)
|
||||
}
|
||||
if got := countPassRows(t, dbURL, "asura"); got != 1 {
|
||||
t.Fatalf("pass rows = %d, want exactly 1", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("asleep", func(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
for i := 0; i < 3; i++ {
|
||||
key := fmt.Sprintf("kagane:w%d", i)
|
||||
seedForCheck(t, s, key, "https://kagane.to/series/"+key[7:], now.Add(-62*time.Minute).UnixMilli())
|
||||
}
|
||||
browser := &fakeFetcher{body: kaganeAPIFixture, status: 200}
|
||||
p := newTestPoller(t, s, &fakeFetcher{status: 200}, now)
|
||||
p.BrowserFetch = browser
|
||||
if pace := p.runLanePass(context.Background(), "kagane", false); pace != defaultGap {
|
||||
t.Fatalf("asleep pace = %s, want %s", pace, defaultGap)
|
||||
}
|
||||
if browser.callCount() != 0 {
|
||||
t.Fatalf("fetches while Chrome is asleep = %d, want 0", browser.callCount())
|
||||
}
|
||||
pass := latestPassFor(t, s, "kagane")
|
||||
if pass.Skip != SkipAsleep || pass.Due != 3 || pass.GapMS != defaultGap.Milliseconds() {
|
||||
t.Fatalf("asleep pass = %+v, want skip %q with 3 due and the default gap", pass, SkipAsleep)
|
||||
}
|
||||
if got := countPassRows(t, dbURL, "kagane"); got != 1 {
|
||||
t.Fatalf("pass rows = %d, want exactly 1", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("eligible-count", func(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
// The eligible query shares the due query's tables, so no real store
|
||||
// failure reaches it after a successful due read; the seam is how the
|
||||
// path is driven at all (see Poller.eligibleCount).
|
||||
p := newTestPoller(t, s, &fakeFetcher{status: 200}, now)
|
||||
p.eligibleCount = func(string) (int, error) { return 0, errors.New("count failed") }
|
||||
if pace := p.runLanePass(context.Background(), "asura", false); pace != defaultGap {
|
||||
t.Fatalf("eligible-count pace = %s, want %s", pace, defaultGap)
|
||||
}
|
||||
if pass := latestPassFor(t, s, "asura"); pass.Skip != SkipEligibleCount {
|
||||
t.Fatalf("skip = %q, want %q", pass.Skip, SkipEligibleCount)
|
||||
}
|
||||
if got := countPassRows(t, dbURL, "asura"); got != 1 {
|
||||
t.Fatalf("pass rows = %d, want exactly 1", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("nothing-eligible", func(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
p := newTestPoller(t, s, &fakeFetcher{status: 200}, now)
|
||||
if pace := p.runLanePass(context.Background(), "asura", false); pace != sites["asura"].Rest {
|
||||
t.Fatalf("nothing-eligible pace = %s, want a full rest %s", pace, sites["asura"].Rest)
|
||||
}
|
||||
if pass := latestPassFor(t, s, "asura"); pass.Skip != SkipNothingEligible {
|
||||
t.Fatalf("skip = %q, want %q", pass.Skip, SkipNothingEligible)
|
||||
}
|
||||
if got := countPassRows(t, dbURL, "asura"); got != 1 {
|
||||
t.Fatalf("pass rows = %d, want exactly 1", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("reached the loop", func(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
const url = "https://asurascans.com/comics/chronicles-of-the-demon-faction-f886a8af"
|
||||
seedForCheck(t, s, "asura:chronicles-of-the-demon-faction-f886a8af", url, 0)
|
||||
p := newTestPoller(t, s, &fakeFetcher{body: asuraSeriesFixture, status: 200}, now)
|
||||
if pace := p.runLanePass(context.Background(), "asura", false); pace != defaultGap {
|
||||
t.Fatalf("loop pace = %s, want %s", pace, defaultGap)
|
||||
}
|
||||
pass := latestPassFor(t, s, "asura")
|
||||
if pass.Skip != "" || pass.Due != 1 || pass.Checked != 1 {
|
||||
t.Fatalf("loop pass = %+v, want an empty skip with 1 due and 1 checked", pass)
|
||||
}
|
||||
if got := countPassRows(t, dbURL, "asura"); got != 1 {
|
||||
t.Fatalf("pass rows = %d, want exactly 1", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("browser-unreachable mid-loop", func(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
seedForCheck(t, s, "comix:c", "https://comix.to/title/c", 0)
|
||||
interrupted := fmt.Errorf("%w: %w", errBrowserInterrupted, errors.New("restart"))
|
||||
browser := &fakeFetcher{status: 200, err: interrupted}
|
||||
p := newTestPoller(t, s, &fakeFetcher{status: 200}, now)
|
||||
p.BrowserFetch = browser
|
||||
if pace := p.runLanePass(context.Background(), "comix", false); pace != defaultGap {
|
||||
t.Fatalf("mid-loop pace = %s, want %s", pace, defaultGap)
|
||||
}
|
||||
// One row was still written, and it is stall-shaped by construction:
|
||||
// empty skip with due > 0 and checked 0 because the return precedes
|
||||
// the checked counter. The assertion below pins that shape precisely.
|
||||
pass := latestPassFor(t, s, "comix")
|
||||
if pass.Skip != "" || pass.Unreachable != 1 || pass.Checked != 0 || pass.Due != 1 {
|
||||
t.Fatalf("mid-loop pass = %+v, want empty skip, unreachable 1, checked 0, due 1", pass)
|
||||
}
|
||||
if got := countPassRows(t, dbURL, "comix"); got != 1 {
|
||||
t.Fatalf("pass rows = %d, want exactly 1", got)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// Carry-forward moves verbatim from the in-memory snapshot (issue #141): a
|
||||
// pass that never computed its own gap carries the previous pass's due, gap,
|
||||
// clamped and checked forward; a pass with a gap of its own records its own
|
||||
// figures, zeroes included, beside the skip reason that explains them.
|
||||
func TestRunLanePassCarryForwardOnlyWhenGapZero(t *testing.T) {
|
||||
now := time.UnixMilli(5_000_000)
|
||||
|
||||
t.Run("no gap of its own carries the previous pass's figures", func(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
seedForCheck(t, s, "kagane:x", "https://kagane.to/series/x", 0)
|
||||
p := newTestPoller(t, s, &fakeFetcher{status: 200}, now)
|
||||
p.BrowserFetch = &fakeFetcher{body: kaganeAPIFixture, status: 200}
|
||||
p.runLanePass(context.Background(), "kagane", false)
|
||||
first := latestPassFor(t, s, "kagane")
|
||||
if first.Due != 1 || first.Checked != 1 || first.GapMS == 0 {
|
||||
t.Fatalf("first pass = %+v, want a measured pass", first)
|
||||
}
|
||||
|
||||
if err := s.SetLaneRefusal("kagane", now.Add(10*time.Minute).UnixMilli()); err != nil {
|
||||
t.Fatalf("SetLaneRefusal: %v", err)
|
||||
}
|
||||
// The pass row is keyed (site, ran_at), so the second pass needs its
|
||||
// own timestamp: a minute later is still inside the refusal.
|
||||
p.Now = func() time.Time { return now.Add(time.Minute) }
|
||||
p.runLanePass(context.Background(), "kagane", false)
|
||||
second := latestPassFor(t, s, "kagane")
|
||||
if second.Skip != SkipRefusing {
|
||||
t.Fatalf("second pass skip = %q, want %q", second.Skip, SkipRefusing)
|
||||
}
|
||||
if second.Due != first.Due || second.Checked != first.Checked ||
|
||||
second.GapMS != first.GapMS || second.Clamped != first.Clamped {
|
||||
t.Fatalf("carried pass = %+v, want the first pass's figures %+v", second, first)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a gap of its own records its own figures", func(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
seedForCheck(t, s, "comix:c", "https://comix.to/title/c", 0)
|
||||
p := newTestPoller(t, s, &fakeFetcher{status: 200}, now)
|
||||
p.BrowserFetch = &fakeFetcher{body: comixSeriesFixture, status: 200}
|
||||
p.runLanePass(context.Background(), "comix", false)
|
||||
first := latestPassFor(t, s, "comix")
|
||||
if first.Due != 1 || first.GapMS == 0 {
|
||||
t.Fatalf("first pass = %+v, want a measured pass", first)
|
||||
}
|
||||
|
||||
// The no-fetcher exit sets its own gap, so no carry-forward: the due
|
||||
// count it never gathered records as zero beside its skip reason. A
|
||||
// minute later gives the second pass its own (site, ran_at) key.
|
||||
p.BrowserFetch = nil
|
||||
p.Now = func() time.Time { return now.Add(time.Minute) }
|
||||
p.runLanePass(context.Background(), "comix", false)
|
||||
second := latestPassFor(t, s, "comix")
|
||||
if second.Skip != SkipNoFetcher {
|
||||
t.Fatalf("second pass skip = %q, want %q", second.Skip, SkipNoFetcher)
|
||||
}
|
||||
if second.GapMS != defaultGap.Milliseconds() {
|
||||
t.Fatalf("second pass gap = %d, want its own %d (not carried)", second.GapMS, defaultGap.Milliseconds())
|
||||
}
|
||||
if second.Due != 0 || second.Checked != 0 {
|
||||
t.Fatalf("second pass = %+v, want due 0 checked 0 (its own, not the previous pass's)", second)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// The pass row's five outcome counts are the classification the Series read
|
||||
// already makes — never a second taxonomy (issue #141). Success is derived,
|
||||
// never stored: checked minus the four named counts, unreachable excluded
|
||||
// because its exit returns before the checked counter increments.
|
||||
func TestRunLanePassCountsOutcomes(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
now := time.UnixMilli(5_000_000)
|
||||
const (
|
||||
successKey = "asura:chronicles-of-the-demon-faction-f886a8af"
|
||||
successURL = "https://asurascans.com/comics/chronicles-of-the-demon-faction-f886a8af"
|
||||
)
|
||||
seeds := map[string]string{
|
||||
"asura:refused": "https://asurascans.com/series/refused",
|
||||
"asura:no-chapter": "https://asurascans.com/comics/no-chapter",
|
||||
"asura:unfetchable": "https://evil.example/x",
|
||||
"asura:transport": "https://asurascans.com/series/transport",
|
||||
successKey: successURL,
|
||||
}
|
||||
for key, url := range seeds {
|
||||
seedForCheck(t, s, key, url, 0)
|
||||
}
|
||||
|
||||
f := &fakeFetcher{perURL: map[string]fakeResponse{
|
||||
seeds["asura:refused"]: {status: 403},
|
||||
seeds["asura:no-chapter"]: {body: "<html></html>", status: 200},
|
||||
seeds["asura:transport"]: {err: errors.New("dial tcp: refused")},
|
||||
successURL: {body: asuraSeriesFixture, status: 200},
|
||||
}}
|
||||
newTestPoller(t, s, f, now).runLanePass(context.Background(), "asura", false)
|
||||
|
||||
pass := latestPassFor(t, s, "asura")
|
||||
if pass.Skip != "" || pass.Due != 5 || pass.Checked != 5 {
|
||||
t.Fatalf("pass = %+v, want a full pass over 5 due Series", pass)
|
||||
}
|
||||
if pass.Refused != 1 || pass.Unreachable != 0 || pass.NoChapter != 1 ||
|
||||
pass.Unfetchable != 1 || pass.Errors != 1 {
|
||||
t.Fatalf("outcome counts = refused %d unreachable %d no_chapter %d unfetchable %d errors %d, want 1 0 1 1 1",
|
||||
pass.Refused, pass.Unreachable, pass.NoChapter, pass.Unfetchable, pass.Errors)
|
||||
}
|
||||
if success := pass.Checked - (pass.Refused + pass.NoChapter + pass.Unfetchable + pass.Errors); success != 1 {
|
||||
t.Fatalf("derived success = %d, want 1", success)
|
||||
}
|
||||
// The one genuine read went through: the success Series carries the
|
||||
// fixture's newest chapter, and none of the four failures do.
|
||||
b, ok, err := s.Get(s.OwnerID(), successKey)
|
||||
if err != nil || !ok {
|
||||
t.Fatalf("Get: %v ok=%v", err, ok)
|
||||
}
|
||||
if b.LatestChapterNum == nil || *b.LatestChapterNum != 181 {
|
||||
t.Fatalf("success Series latest = %v, want 181 (the fixture's newest)", b.LatestChapterNum)
|
||||
}
|
||||
for _, key := range []string{"asura:refused", "asura:no-chapter", "asura:unfetchable", "asura:transport"} {
|
||||
b, _, err := s.Get(s.OwnerID(), key)
|
||||
if err != nil {
|
||||
t.Fatalf("Get %s: %v", key, err)
|
||||
}
|
||||
if b.LatestChapterNum != nil {
|
||||
t.Fatalf("%s latest = %v, want nil (no chapter survived a failed read)", key, *b.LatestChapterNum)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// A refusal is the Site's mood and outlives our process: the durable stamp a
|
||||
// pass writes is honoured by a freshly constructed poller, which must not
|
||||
// re-probe the Site inside its backoff (issue #141).
|
||||
func TestDurableRefusalSurvivesFreshPoller(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
now := time.UnixMilli(5_000_000)
|
||||
// Four due Series: the Lane refuses twice, stamps those two, and leaves
|
||||
// the remaining two untried and due — the queue the fresh poller must
|
||||
// find intact once the durable backoff lifts.
|
||||
for i := 0; i < 4; i++ {
|
||||
key := fmt.Sprintf("kagane:s%d", i)
|
||||
seedForCheck(t, s, key, "https://kagane.to/series/"+key[7:], 0)
|
||||
}
|
||||
|
||||
p := newTestPoller(t, s, &fakeFetcher{status: 200}, now)
|
||||
p.BrowserFetch = &fakeFetcher{status: 403}
|
||||
p.runLanePass(context.Background(), "kagane", false)
|
||||
if got := p.BrowserFetch.(*fakeFetcher).callCount(); got != 2 {
|
||||
t.Fatalf("fetches on the refusing pass = %d, want 2 (refused twice)", got)
|
||||
}
|
||||
|
||||
// A restart: a brand-new poller, no in-memory refusal, same store. The
|
||||
// durable stamp gates the pass.
|
||||
fresh := newTestPoller(t, s, &fakeFetcher{status: 200}, now.Add(14*time.Minute))
|
||||
fresh.BrowserFetch = &fakeFetcher{status: 403}
|
||||
fresh.runLanePass(context.Background(), "kagane", false)
|
||||
if got := fresh.BrowserFetch.(*fakeFetcher).callCount(); got != 0 {
|
||||
t.Fatalf("fetches by a fresh poller inside the backoff = %d, want 0", got)
|
||||
}
|
||||
if pass := latestPassFor(t, s, "kagane"); pass.Skip != SkipRefusing {
|
||||
t.Fatalf("fresh poller's pass skip = %q, want %q", pass.Skip, SkipRefusing)
|
||||
}
|
||||
|
||||
// Past the backoff the fresh poller probes again — the two Series the
|
||||
// original pass never reached, still due with their stamps untouched.
|
||||
fresh.Now = func() time.Time { return now.Add(16 * time.Minute) }
|
||||
fresh.runLanePass(context.Background(), "kagane", false)
|
||||
if got := fresh.BrowserFetch.(*fakeFetcher).callCount(); got != 2 {
|
||||
t.Fatalf("fetches after the backoff = %d, want 2", got)
|
||||
}
|
||||
for i := 2; i < 4; i++ {
|
||||
if got := readLatestCheckedAt(t, s, fmt.Sprintf("kagane:s%d", i)); got == 0 {
|
||||
t.Fatalf("kagane:s%d still untried after the backoff", i)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The pause is read ahead of the refusal check: a Lane that is both paused
|
||||
// and inside a refusal backoff records the paused skip value, not the
|
||||
// refusing one. The pause is the owner's order and outranks the Site's mood
|
||||
// (issue #147).
|
||||
func TestPauseGatePrecedesRefusalGate(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
now := time.UnixMilli(5_000_000)
|
||||
seedForCheck(t, s, "asura:x", "https://asurascans.com/series/x", 0)
|
||||
if err := s.SetLaneRefusal("asura", now.Add(10*time.Minute).UnixMilli()); err != nil {
|
||||
t.Fatalf("SetLaneRefusal: %v", err)
|
||||
}
|
||||
if err := s.PauseLane("asura", now.Add(30*time.Minute).UnixMilli()); err != nil {
|
||||
t.Fatalf("PauseLane: %v", err)
|
||||
}
|
||||
|
||||
f := &fakeFetcher{body: asuraSeriesFixture, status: 200}
|
||||
p := newTestPoller(t, s, f, now)
|
||||
if pace := p.runLanePass(context.Background(), "asura", false); pace != 30*time.Minute {
|
||||
t.Fatalf("paused-while-refusing pace = %s, want 30m (the pause's expiry)", pace)
|
||||
}
|
||||
if f.callCount() != 0 {
|
||||
t.Fatalf("fetches while paused and refusing = %d, want 0", f.callCount())
|
||||
}
|
||||
if pass := latestPassFor(t, s, "asura"); pass.Skip != SkipPaused {
|
||||
t.Fatalf("skip = %q, want %q (the pause outranks the refusal)", pass.Skip, SkipPaused)
|
||||
}
|
||||
}
|
||||
|
||||
// A pause is a fact about the Site, not about the process: a freshly
|
||||
// constructed poller against a store holding a pause row stays paused until
|
||||
// the expiry, then runs the Lane normally. The restart criterion is the
|
||||
// whole point of writing a row instead of commanding a poller (issue #147).
|
||||
func TestDurablePauseSurvivesFreshPoller(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
now := time.UnixMilli(5_000_000)
|
||||
seedForCheck(t, s, "asura:x", "https://asurascans.com/series/x", 0)
|
||||
if err := s.PauseLane("asura", now.Add(30*time.Minute).UnixMilli()); err != nil {
|
||||
t.Fatalf("PauseLane: %v", err)
|
||||
}
|
||||
|
||||
// A restart: a brand-new poller, no in-memory state, same store.
|
||||
fresh := newTestPoller(t, s, &fakeFetcher{body: asuraSeriesFixture, status: 200}, now)
|
||||
if pace := fresh.runLanePass(context.Background(), "asura", false); pace != 30*time.Minute {
|
||||
t.Fatalf("fresh poller's paused pace = %s, want 30m", pace)
|
||||
}
|
||||
if pass := latestPassFor(t, s, "asura"); pass.Skip != SkipPaused {
|
||||
t.Fatalf("fresh poller's pass skip = %q, want %q", pass.Skip, SkipPaused)
|
||||
}
|
||||
|
||||
// Past the expiry the same fresh poller runs the Lane normally.
|
||||
fresh.Now = func() time.Time { return now.Add(31 * time.Minute) }
|
||||
fresh.runLanePass(context.Background(), "asura", false)
|
||||
if pass := latestPassFor(t, s, "asura"); pass.Skip != "" {
|
||||
t.Fatalf("pass after the expiry skip = %q, want the loop reached", pass.Skip)
|
||||
}
|
||||
if got := readLatestCheckedAt(t, s, "asura:x"); got == 0 {
|
||||
t.Fatal("the Series was not checked after the pause lifted")
|
||||
}
|
||||
}
|
||||
|
||||
// ResumeLane zeroes the pause and the Lane's next pass finds its full queue
|
||||
// waiting: a pause delays work rather than discarding it, so the due Series
|
||||
// sit unstamped while paused and are all fetched once the pause lifts
|
||||
// (issue #147).
|
||||
func TestResumeLaneRestoresTheQueue(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
now := time.UnixMilli(5_000_000)
|
||||
for i := 0; i < 3; i++ {
|
||||
seedForCheck(t, s, fmt.Sprintf("asura:s%d", i), "https://asurascans.com/series/x", 0)
|
||||
}
|
||||
if err := s.PauseLane("asura", now.Add(30*time.Minute).UnixMilli()); err != nil {
|
||||
t.Fatalf("PauseLane: %v", err)
|
||||
}
|
||||
|
||||
f := &fakeFetcher{body: asuraSeriesFixture, status: 200}
|
||||
p := newTestPoller(t, s, f, now)
|
||||
p.runLanePass(context.Background(), "asura", false)
|
||||
if f.callCount() != 0 {
|
||||
t.Fatalf("fetches while paused = %d, want 0", f.callCount())
|
||||
}
|
||||
for i := 0; i < 3; i++ {
|
||||
if got := readLatestCheckedAt(t, s, fmt.Sprintf("asura:s%d", i)); got != 0 {
|
||||
t.Fatalf("asura:s%d stamp while paused = %d, want 0 (due and unstamped)", i, got)
|
||||
}
|
||||
}
|
||||
|
||||
if err := s.ResumeLane("asura"); err != nil {
|
||||
t.Fatalf("ResumeLane: %v", err)
|
||||
}
|
||||
before := f.callCount()
|
||||
p.runLanePass(context.Background(), "asura", false)
|
||||
if got := f.callCount() - before; got != 3 {
|
||||
t.Fatalf("fetches after resume = %d, want 3 (the full due queue)", got)
|
||||
}
|
||||
for i := 0; i < 3; i++ {
|
||||
if got := readLatestCheckedAt(t, s, fmt.Sprintf("asura:s%d", i)); got == 0 {
|
||||
t.Fatalf("asura:s%d still untried after resume", i)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// RecordLanePass prunes in the same call that inserts, so the retention
|
||||
// cutoff the recorder passes is observable in what survives: a row just inside
|
||||
// 14 days behind the poller's clock is kept, one just outside is pruned
|
||||
// (issue #139, #141).
|
||||
func TestPassRetentionCutoffIsFourteenDays(t *testing.T) {
|
||||
s, dbURL := newTestStore(t)
|
||||
// A real-world clock: the seeded rows sit 14 days back, so they must be
|
||||
// positive timestamps or the seed's own prune (ran_at < 0) removes them.
|
||||
now := time.UnixMilli(1_800_000_000_000)
|
||||
kept := now.Add(-14*24*time.Hour + time.Minute).UnixMilli()
|
||||
pruned := now.Add(-14*24*time.Hour - time.Minute).UnixMilli()
|
||||
for _, ranAt := range []int64{kept, pruned} {
|
||||
if err := s.RecordLanePass(store.LanePass{Site: "asura", RanAt: ranAt}, 0); err != nil {
|
||||
t.Fatalf("seed pass at %d: %v", ranAt, err)
|
||||
}
|
||||
}
|
||||
|
||||
seedForCheck(t, s, "asura:x", "https://asurascans.com/series/x", 0)
|
||||
newTestPoller(t, s, &fakeFetcher{body: asuraSeriesFixture, status: 200}, now).
|
||||
runLanePass(context.Background(), "asura", false)
|
||||
|
||||
db, err := sql.Open("pgx", dbURL)
|
||||
if err != nil {
|
||||
t.Fatalf("open %s: %v", dbURL, err)
|
||||
}
|
||||
defer db.Close()
|
||||
var n int
|
||||
if err := db.QueryRow(`SELECT count(*) FROM poll_passes WHERE site = $1 AND ran_at = $2`, "asura", pruned).Scan(&n); err != nil {
|
||||
t.Fatalf("count pruned row: %v", err)
|
||||
}
|
||||
if n != 0 {
|
||||
t.Fatalf("row at %d survived, want it pruned (older than 14 days)", pruned)
|
||||
}
|
||||
if err := db.QueryRow(`SELECT count(*) FROM poll_passes WHERE site = $1`, "asura").Scan(&n); err != nil {
|
||||
t.Fatalf("count passes: %v", err)
|
||||
}
|
||||
if n != 2 {
|
||||
t.Fatalf("passes = %d, want 2 (this pass plus the kept row)", n)
|
||||
}
|
||||
}
|
||||
|
||||
// stampRecordingFetcher is a fakeFetcher that records the Series' check stamp
|
||||
// at call time — the seam that proves the check stamp is written before the
|
||||
// fetch (issue #146). If the order were swapped, the recorded stamp would be
|
||||
// the pre-pass value and the ordering assertion would fail.
|
||||
type stampRecordingFetcher struct {
|
||||
fakeFetcher
|
||||
store *store.Store
|
||||
site string
|
||||
seriesID string
|
||||
stampAtCall int64
|
||||
}
|
||||
|
||||
func (f *stampRecordingFetcher) Get(ctx context.Context, url string) (string, int, error) {
|
||||
ts, err := f.store.LatestCheckedAt(f.site, f.seriesID)
|
||||
if err != nil {
|
||||
return "", 0, err
|
||||
}
|
||||
f.stampAtCall = ts
|
||||
return f.fakeFetcher.Get(ctx, url)
|
||||
}
|
||||
|
||||
// The check stamp is written before the fetch is attempted: the fake fetcher
|
||||
// records the stamp it sees at call time, and it must already be the pass's
|
||||
// own stamp. The order is load-bearing — a forced request is pending while
|
||||
// force_poll_at > latest_checked_at, so stamping after the fetch would make a
|
||||
// failed forced request sticky — and the assertion fails if it is swapped.
|
||||
func TestCheckStampIsWrittenBeforeFetch(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
const (
|
||||
key = "asura:chronicles-of-the-demon-faction-f886a8af"
|
||||
seriesURL = "https://asurascans.com/comics/chronicles-of-the-demon-faction-f886a8af"
|
||||
)
|
||||
seedForCheck(t, s, key, seriesURL, 0)
|
||||
now := time.UnixMilli(7_000_000)
|
||||
rec := &stampRecordingFetcher{
|
||||
fakeFetcher: fakeFetcher{body: asuraSeriesFixture, status: 200},
|
||||
store: s, site: "asura", seriesID: "chronicles-of-the-demon-faction-f886a8af",
|
||||
}
|
||||
newTestPoller(t, s, rec, now).runOnce(context.Background())
|
||||
if rec.stampAtCall != now.UnixMilli() {
|
||||
t.Fatalf("check stamp at fetch time = %d, want %d (the stamp must be written before the fetch)", rec.stampAtCall, now.UnixMilli())
|
||||
}
|
||||
}
|
||||
|
||||
// A forced Series whose attempt fails is no longer pending: the check stamp
|
||||
// is written before the fetch, so the first attempt ends the pending state
|
||||
// whatever it returns. A naive implementation (stamp only on success, or
|
||||
// after the fetch) leaves the request sticky and the row due next pass.
|
||||
func TestForcedSeriesSelfClearsOnFailedAttempt(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
now := time.UnixMilli(7_000_000)
|
||||
const url = "https://asurascans.com/comics/x"
|
||||
// Freshly checked, so only the force flag makes it due.
|
||||
seedForCheck(t, s, "asura:x", url, now.Add(-30*time.Minute).UnixMilli())
|
||||
if err := s.ForceSeriesPoll("asura", "x", now.UnixMilli()); err != nil {
|
||||
t.Fatalf("ForceSeriesPoll: %v", err)
|
||||
}
|
||||
|
||||
due, err := s.DueForLatestCheck("asura", now.Add(-time.Hour).UnixMilli(), -1)
|
||||
if err != nil {
|
||||
t.Fatalf("DueForLatestCheck: %v", err)
|
||||
}
|
||||
if len(due) != 1 || due[0].Key() != "asura:x" || !due[0].Forced {
|
||||
t.Fatalf("forced row not due and flagged before the attempt: %+v", due)
|
||||
}
|
||||
|
||||
// The attempt fails, but the attempt still happened: the row is stamped
|
||||
// and no longer pending.
|
||||
f := &fakeFetcher{err: errors.New("dial tcp: refused")}
|
||||
newTestPoller(t, s, f, now).runOnce(context.Background())
|
||||
if got := readLatestCheckedAt(t, s, "asura:x"); got != now.UnixMilli() {
|
||||
t.Fatalf("latest_checked_at = %d, want %d", got, now.UnixMilli())
|
||||
}
|
||||
due, err = s.DueForLatestCheck("asura", now.Add(-time.Hour).UnixMilli(), -1)
|
||||
if err != nil {
|
||||
t.Fatalf("DueForLatestCheck: %v", err)
|
||||
}
|
||||
if len(due) != 0 {
|
||||
t.Fatalf("failed attempt left the forced row due: %v", due)
|
||||
}
|
||||
}
|
||||
|
||||
// The refusal backoff is a Lane gate, not a Series gate: a forced Series does
|
||||
// not override a Site that is actively refusing, because hand-forcing a
|
||||
// request into a refusal only makes it worse (issue #146).
|
||||
func TestForcedSeriesDoesNotOverrideRefusalBackoff(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
now := time.UnixMilli(5_000_000)
|
||||
seedForCheck(t, s, "kagane:w", "https://kagane.to/series/w", 0)
|
||||
if err := s.ForceSeriesPoll("kagane", "w", now.UnixMilli()); err != nil {
|
||||
t.Fatalf("ForceSeriesPoll: %v", err)
|
||||
}
|
||||
if err := s.SetLaneRefusal("kagane", now.Add(RefuseBackoff).UnixMilli()); err != nil {
|
||||
t.Fatalf("SetLaneRefusal: %v", err)
|
||||
}
|
||||
|
||||
browser := &fakeFetcher{body: kaganeAPIFixture, status: 200}
|
||||
p := &Poller{
|
||||
Store: s, Fetch: &fakeFetcher{status: 200}, BrowserFetch: browser,
|
||||
Now: func() time.Time { return now },
|
||||
}
|
||||
p.runOnce(context.Background())
|
||||
if got := browser.callCount(); got != 0 {
|
||||
t.Fatalf("browser fetches through a refusal backoff = %d, want 0", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A forced Series wakes a sleeping browser Lane: the wake thresholds exist to
|
||||
// stop the machine waking itself for one unattended check, and a human asking
|
||||
// is not that (issue #146). Below both thresholds the Lane still sleeps when
|
||||
// nothing is forced.
|
||||
func TestForcedSeriesWakesSleepingBrowser(t *testing.T) {
|
||||
s, _ := newTestStore(t)
|
||||
now := time.UnixMilli(5_000_000)
|
||||
// Freshly due: checked two minutes before the rest elapses, so the wait
|
||||
// is far below browserWakeAge and the count is under browserWakeCount.
|
||||
seedForCheck(t, s, "kagane:w1", "https://kagane.to/series/w1", now.Add(-62*time.Minute).UnixMilli())
|
||||
seedForCheck(t, s, "kagane:w2", "https://kagane.to/series/w2", now.Add(-62*time.Minute).UnixMilli())
|
||||
|
||||
browser := &fakeFetcher{body: kaganeAPIFixture, status: 200}
|
||||
p := &Poller{
|
||||
Store: s, Fetch: &fakeFetcher{status: 200}, BrowserFetch: browser,
|
||||
Now: func() time.Time { return now },
|
||||
}
|
||||
p.runOnce(context.Background())
|
||||
if got := browser.callCount(); got != 0 {
|
||||
t.Fatalf("browser fetches without a forced series = %d, want 0 (Chrome stays asleep)", got)
|
||||
}
|
||||
|
||||
if err := s.ForceSeriesPoll("kagane", "w1", now.UnixMilli()); err != nil {
|
||||
t.Fatalf("ForceSeriesPoll: %v", err)
|
||||
}
|
||||
p.runOnce(context.Background())
|
||||
if got := browser.callCount(); got != 2 {
|
||||
t.Fatalf("browser fetches with a forced series = %d, want 2 (the lane wakes)", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -401,9 +401,10 @@ const (
|
||||
// express (docs/research/cloudflare-bot-scoring-and-poll-cadence.md);
|
||||
// below it the Lane is outrunning its own plan and says so loudly.
|
||||
minGap = time.Second
|
||||
// refuseBackoff is how long a Lane waits after its Site refused twice in
|
||||
// one run before attempting it again.
|
||||
refuseBackoff = 15 * time.Minute
|
||||
// RefuseBackoff is how long a Lane waits after its Site refused twice in
|
||||
// one run before attempting it again. Exported so the web layer can derive
|
||||
// browser reachability from the pass log over the same window (issue #145).
|
||||
RefuseBackoff = 15 * time.Minute
|
||||
// browserWakeCount and browserWakeAge gate a browser Lane's run: five or
|
||||
// more due Series, or any one of them waiting this long, or Chrome stays
|
||||
// asleep (ADR-0005 on-demand browser).
|
||||
@@ -511,6 +512,18 @@ var sites = map[string]site{
|
||||
},
|
||||
}
|
||||
|
||||
// SiteNames returns every registry Site, sorted. The admin Series list's Site
|
||||
// select needs the full registry, not just the Sites that have rows, and
|
||||
// laneNames() is the poller's copy of the same list — both read this.
|
||||
func SiteNames() []string {
|
||||
names := make([]string, 0, len(sites))
|
||||
for name := range sites {
|
||||
names = append(names, name)
|
||||
}
|
||||
sort.Strings(names)
|
||||
return names
|
||||
}
|
||||
|
||||
// browserBackedSites is derived from the registry: the Sites whose pages are
|
||||
// read through the browser sidecar. Sorted so callers that range it (the
|
||||
// browser fetcher's dispatch) see a stable order instead of map-iteration
|
||||
|
||||
@@ -1,79 +0,0 @@
|
||||
package latest
|
||||
|
||||
import "time"
|
||||
|
||||
// LaneState is the administrative page's view of one Poll Lane (issue #102):
|
||||
// what the Lane's last pass saw. Due, Gap and Checked are filled in as the
|
||||
// pass computes them; a pass that returned before reaching a figure (refusal
|
||||
// backoff, sidecar down) carries the previous pass's figures forward rather
|
||||
// than overwriting them with zeroes the page would state as fact.
|
||||
type LaneState struct {
|
||||
Site string
|
||||
Due int
|
||||
LastRun time.Time
|
||||
Gap time.Duration
|
||||
// Checked is how many Series this pass actually read. A Lane with Series
|
||||
// due and nothing checked has stopped working; one with nothing due is
|
||||
// merely quiet, and the page must not draw the two the same (story 13).
|
||||
Checked int
|
||||
Clamped bool
|
||||
Refusing bool
|
||||
Browser bool
|
||||
// Asleep marks a browser Lane whose last pass declined to wake Chrome
|
||||
// because it was under both wake thresholds (ADR-0005). Due without
|
||||
// Checked then means "waiting for the group to gather", not "stopped", and
|
||||
// the page must not draw it as a stall.
|
||||
Asleep bool
|
||||
}
|
||||
|
||||
// Status is the owner's page snapshot of the whole poller (issue #102).
|
||||
type Status struct {
|
||||
Lanes []LaneState
|
||||
BrowserConfigured bool
|
||||
BrowserReachable bool
|
||||
}
|
||||
|
||||
// LaneStatus returns a copy of the poller's Lane state for the owner's page.
|
||||
// Only Sites that have completed a pass appear — a restart therefore renders
|
||||
// "no data yet" instead of confident zeroes — in the same order Run iterates.
|
||||
// Refusing is derived at snapshot time from the refusal backoff, not stored,
|
||||
// so a Lane that cooled down between passes reports false without a new pass.
|
||||
// BrowserReachable mirrors the Lanes' own gate: the sidecar is down only
|
||||
// within the refuseBackoff window since its last loss.
|
||||
func (p *Poller) LaneStatus() Status {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
lanes := make([]LaneState, 0, len(p.laneStates))
|
||||
now := p.Now()
|
||||
for _, name := range laneNames() {
|
||||
st, ok := p.laneStates[name]
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
st.Refusing = now.Before(p.refuseUntil[name])
|
||||
lanes = append(lanes, st)
|
||||
}
|
||||
configured := p.BrowserFetch != nil
|
||||
reachable := configured
|
||||
if reachable && !p.browserDownAt.IsZero() && now.Sub(p.browserDownAt) < refuseBackoff {
|
||||
reachable = false
|
||||
}
|
||||
return Status{Lanes: lanes, BrowserConfigured: configured, BrowserReachable: reachable}
|
||||
}
|
||||
|
||||
// recordLaneState stores one Lane's last pass for LaneStatus. Called deferred
|
||||
// from runLanePass so every return path records, even a pass that refused.
|
||||
// A pass that never reached the pace (Gap zero) keeps the last pass's figures:
|
||||
// the Lane's due count and gap did not become zero because this pass declined
|
||||
// to look, and the row's own marks say why it declined.
|
||||
func (p *Poller) recordLaneState(st LaneState) {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
if p.laneStates == nil {
|
||||
p.laneStates = make(map[string]LaneState)
|
||||
}
|
||||
if prev, ok := p.laneStates[st.Site]; ok && st.Gap == 0 {
|
||||
st.Due, st.Gap, st.Clamped, st.Checked = prev.Due, prev.Gap, prev.Clamped, prev.Checked
|
||||
}
|
||||
p.laneStates[st.Site] = st
|
||||
}
|
||||
Reference in New Issue
Block a user