feat: one shared series-page read for the Poll and the Acquisition (#94)
Extract gate, route, fetch and both parses into readSeriesPage, which returns facts only; checkOne and acquire keep their scheduling, stamp order and persistence policies. The poller also stops logging the fetched body length, which only made sense while the body lived at the call site.
This commit is contained in:
@@ -2,6 +2,7 @@ package latest
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
@@ -97,51 +98,49 @@ func (a *Acquirer) acquire(ctx context.Context, sr store.Series) {
|
||||
if a.Fetch == nil && a.BrowserFetch == nil {
|
||||
return
|
||||
}
|
||||
// series_url arrives in a client-supplied PUT body, so the same gate the
|
||||
// poller uses applies here — without it a token-holder chooses what the
|
||||
// server fetches from its own network position.
|
||||
if !fetchableSeriesURL(sr.Site, sr.SeriesURL) {
|
||||
log.Printf("acquire %q: not fetchable: site=%q url=%q", sr.Key(), sr.Site, sr.SeriesURL)
|
||||
return
|
||||
}
|
||||
|
||||
f := fetcherFor(sr.Site, a.BrowserFetch, a.Fetch)
|
||||
if f == nil {
|
||||
log.Printf("acquire %q: no fetcher for site %q", sr.Key(), sr.Site)
|
||||
return
|
||||
}
|
||||
body, status, err := f.Get(ctx, sr.SeriesURL)
|
||||
facts, err := readSeriesPage(ctx, sr.Site, sr.SeriesURL, a.BrowserFetch, a.Fetch)
|
||||
if err != nil {
|
||||
log.Printf("acquire %q: fetch %s: %v", sr.Key(), sr.SeriesURL, err)
|
||||
return
|
||||
switch {
|
||||
case errors.Is(err, errNotFetchable):
|
||||
// series_url arrives in a client-supplied PUT body, so without the
|
||||
// gate a token-holder chooses what the server fetches from its own
|
||||
// network position.
|
||||
log.Printf("acquire %q: not fetchable: site=%q url=%q", sr.Key(), sr.Site, sr.SeriesURL)
|
||||
case errors.Is(err, errNoFetcher):
|
||||
log.Printf("acquire %q: no fetcher for site %q", sr.Key(), sr.Site)
|
||||
default:
|
||||
log.Printf("acquire %q: %v", sr.Key(), err)
|
||||
}
|
||||
if status != 200 {
|
||||
log.Printf("acquire %q: fetch %s: status %d", sr.Key(), sr.SeriesURL, status)
|
||||
return
|
||||
}
|
||||
|
||||
// This page just served the same purpose a poll tick would have; without
|
||||
// the stamp the row stays due and the poller refetches it immediately.
|
||||
//
|
||||
// Stamped after success — the reverse of the poller, which stamps before
|
||||
// the fetch: the Reader is here, watching the Series they just created, so
|
||||
// a failed acquisition must leave the row due for a fast retry rather than
|
||||
// consuming the cooldown. The stamp happens even when the page read
|
||||
// succeeded but produced no facts to persist.
|
||||
if err := a.Store.MarkLatestChecked(sr.Site, sr.SeriesID, time.Now().UnixMilli()); err != nil {
|
||||
log.Printf("acquire %q: mark checked: %v", sr.Key(), err)
|
||||
}
|
||||
|
||||
if latest, ok := latestChapterFrom(sr.Site, sr.SeriesURL, body); ok {
|
||||
if err := a.Store.SetLatestChapter(sr.Site, sr.SeriesID, latest.Label, latest.Num); err != nil {
|
||||
if facts.HasLatest {
|
||||
if err := a.Store.SetLatestChapter(sr.Site, sr.SeriesID, facts.Latest.Label, facts.Latest.Num); err != nil {
|
||||
log.Printf("acquire %q: set latest chapter: %v", sr.Key(), err)
|
||||
}
|
||||
}
|
||||
|
||||
cover, ok := coverFrom(sr.Site, sr.SeriesURL, body)
|
||||
if !ok {
|
||||
if !facts.HasCover {
|
||||
return
|
||||
}
|
||||
bytes, contentType, err := fetchCoverBytes(ctx, cover, a.BrowserCoverFetch, a.Covers)
|
||||
bytes, contentType, err := fetchCoverBytes(ctx, facts.Cover, a.BrowserCoverFetch, a.Covers)
|
||||
if err != nil {
|
||||
log.Printf("acquire %q: fetch cover %s: %v", sr.Key(), cover, err)
|
||||
log.Printf("acquire %q: fetch cover %s: %v", sr.Key(), facts.Cover, err)
|
||||
return
|
||||
}
|
||||
if err := a.Store.SetSeriesCover(sr.Site, sr.SeriesID, cover, bytes, contentType); err != nil {
|
||||
if err := a.Store.SetSeriesCover(sr.Site, sr.SeriesID, facts.Cover, bytes, contentType); err != nil {
|
||||
log.Printf("acquire %q: persist cover: %v", sr.Key(), err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ package latest
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log"
|
||||
"net/url"
|
||||
"time"
|
||||
@@ -62,17 +63,17 @@ type Poller struct {
|
||||
// failed, both heal here. A non-blank CoverAddress is left alone — refetching
|
||||
// would add a request per Series per cycle and change artwork under the Reader
|
||||
// for no visible reason. A row that already carries a source URL is owned by
|
||||
// prefetchCover instead; this path only extracts from the series page.
|
||||
// prefetchCover instead; this path only records a Cover address already
|
||||
// extracted from the series page.
|
||||
//
|
||||
// Failures are logged against the Series and never returned: the chapter poll
|
||||
// must not notice. A failed fill is retried the next time this Series is due;
|
||||
// there is no separate retry queue.
|
||||
func (p *Poller) fillBlankCover(ctx context.Context, sr store.Series, body string) {
|
||||
func (p *Poller) fillBlankCover(ctx context.Context, sr store.Series, cover string) {
|
||||
if sr.CoverAddress != "" || sr.Cover != "" {
|
||||
return
|
||||
}
|
||||
cover, ok := coverFrom(sr.Site, sr.SeriesURL, body)
|
||||
if !ok {
|
||||
if cover == "" {
|
||||
return
|
||||
}
|
||||
p.storeCover(ctx, sr, cover)
|
||||
@@ -217,43 +218,30 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) {
|
||||
return
|
||||
}
|
||||
|
||||
// series_url is client-supplied (PUT /bookmarks/{key} accepts any string),
|
||||
// so this is not just an optimisation against burning a request on an
|
||||
// unknown site: without it, the server would issue a GET from its own
|
||||
// network position to whatever URL a token-holder writes, including
|
||||
// link-local/internal addresses or non-https schemes. The cooldown above
|
||||
// is already consumed, so a row that never passes this check is retried at
|
||||
// cooldown pace rather than hot-looping.
|
||||
if !fetchableSeriesURL(sr.Site, sr.SeriesURL) {
|
||||
log.Printf("latest poll %q: not fetchable: site=%q url=%q", sr.Key(), sr.Site, sr.SeriesURL)
|
||||
return
|
||||
}
|
||||
|
||||
f := fetcherFor(sr.Site, p.BrowserFetch, p.Fetch)
|
||||
if f == nil {
|
||||
log.Printf("latest poll %q: no fetcher for site %q", sr.Key(), sr.Site)
|
||||
return
|
||||
}
|
||||
p.prefetchCover(ctx, sr)
|
||||
|
||||
body, status, err := f.Get(ctx, sr.SeriesURL)
|
||||
facts, err := readSeriesPage(ctx, sr.Site, sr.SeriesURL, p.BrowserFetch, p.Fetch)
|
||||
if err != nil {
|
||||
log.Printf("latest poll %q: fetch %s: %v", sr.Key(), sr.SeriesURL, err)
|
||||
switch {
|
||||
case errors.Is(err, errNotFetchable):
|
||||
// The cooldown above is already consumed, so a row that never
|
||||
// passes the gate is retried at cooldown pace rather than
|
||||
// hot-looping.
|
||||
log.Printf("latest poll %q: not fetchable: site=%q url=%q", sr.Key(), sr.Site, sr.SeriesURL)
|
||||
case errors.Is(err, errNoFetcher):
|
||||
log.Printf("latest poll %q: no fetcher for site %q", sr.Key(), sr.Site)
|
||||
default:
|
||||
log.Printf("latest poll %q: %v", sr.Key(), err)
|
||||
}
|
||||
return
|
||||
}
|
||||
if status != 200 {
|
||||
log.Printf("latest poll %q: fetch %s: status %d", sr.Key(), sr.SeriesURL, status)
|
||||
return
|
||||
}
|
||||
|
||||
latest, ok := latestChapterFrom(sr.Site, sr.SeriesURL, body)
|
||||
// A legacy cover source is healed independently of the page read.
|
||||
p.prefetchCover(ctx, sr)
|
||||
// Cover fill is independent of the chapter signal: a page that lost its
|
||||
// chapter list may keep its og:image, and a blank Series heals either way.
|
||||
p.fillBlankCover(ctx, sr, body)
|
||||
if !ok {
|
||||
p.fillBlankCover(ctx, sr, facts.Cover)
|
||||
if !facts.HasLatest {
|
||||
// Most likely a challenge page or a layout change. Either way the row is
|
||||
// already stamped, so this waits out a cooldown instead of hot-looping.
|
||||
log.Printf("latest poll %q: no chapter links in %d bytes", sr.Key(), len(body))
|
||||
log.Printf("latest poll %q: no chapter links", sr.Key())
|
||||
return
|
||||
}
|
||||
|
||||
@@ -261,7 +249,7 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) {
|
||||
// chapter should correct the stored number downward. The comparison is
|
||||
// against the due-query snapshot; a concurrent write in between only costs
|
||||
// one redundant UPDATE of the same absolute value, never a wrong one.
|
||||
if sr.LatestChapterNum != nil && *sr.LatestChapterNum == latest.Num {
|
||||
if sr.LatestChapterNum != nil && *sr.LatestChapterNum == facts.Latest.Num {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -269,11 +257,11 @@ func (p *Poller) checkOne(ctx context.Context, sr store.Series) {
|
||||
// bookmark joining to it, and the bookmark's updated_at is never touched —
|
||||
// a newly published chapter is not reading progress and must not reorder
|
||||
// the list.
|
||||
if err := p.Store.SetLatestChapter(sr.Site, sr.SeriesID, latest.Label, latest.Num); err != nil {
|
||||
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
|
||||
}
|
||||
log.Printf("latest poll %q: latest is now %s", sr.Key(), latest.Label)
|
||||
log.Printf("latest poll %q: latest is now %s", sr.Key(), facts.Latest.Label)
|
||||
}
|
||||
|
||||
// fetchableSeriesURL reports whether site is a Site the registry knows and
|
||||
|
||||
@@ -0,0 +1,55 @@
|
||||
package latest
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
)
|
||||
|
||||
// seriesRead carries the two facts the poll and the acquirer both extract
|
||||
// from a series page. Persistence, stamps and scheduling stay with the
|
||||
// callers, so the policies that keep the two flows distinct (stamp order,
|
||||
// cooldowns) are not swallowed by the module.
|
||||
type seriesRead struct {
|
||||
Latest latestChapter
|
||||
HasLatest bool
|
||||
Cover string
|
||||
HasCover bool
|
||||
}
|
||||
|
||||
// errNotFetchable and errNoFetcher separate the gate and the route from fetch
|
||||
// failures so each caller keeps its own distinct log line for all three.
|
||||
var (
|
||||
errNotFetchable = errors.New("series url not fetchable")
|
||||
errNoFetcher = errors.New("no fetcher for site")
|
||||
)
|
||||
|
||||
// readSeriesPage performs the series-page read the poll and the acquirer have
|
||||
// in common: gate the address, choose the route, fetch the page, extract the
|
||||
// Latest Chapter and the Cover address. It persists nothing and stamps
|
||||
// nothing.
|
||||
//
|
||||
// series_url arrives in a client-supplied PUT body (PUT /bookmarks/{key}
|
||||
// accepts any string), so the gate is not an optimisation against burning a
|
||||
// request on an unknown site: without it, the server would issue a GET from
|
||||
// its own network position to whatever URL a token-holder writes, including
|
||||
// link-local/internal addresses or non-https schemes.
|
||||
func readSeriesPage(ctx context.Context, site, seriesURL string, browser, tls Fetcher) (seriesRead, error) {
|
||||
if !fetchableSeriesURL(site, seriesURL) {
|
||||
return seriesRead{}, fmt.Errorf("%w: site=%q url=%q", errNotFetchable, site, seriesURL)
|
||||
}
|
||||
f := fetcherFor(site, browser, tls)
|
||||
if f == nil {
|
||||
return seriesRead{}, fmt.Errorf("%w: site %q", errNoFetcher, site)
|
||||
}
|
||||
body, status, err := f.Get(ctx, seriesURL)
|
||||
if err != nil {
|
||||
return seriesRead{}, fmt.Errorf("fetch %s: %w", seriesURL, err)
|
||||
}
|
||||
if status != 200 {
|
||||
return seriesRead{}, fmt.Errorf("fetch %s: status %d", seriesURL, status)
|
||||
}
|
||||
latest, hasLatest := latestChapterFrom(site, seriesURL, body)
|
||||
cover, hasCover := coverFrom(site, seriesURL, body)
|
||||
return seriesRead{Latest: latest, HasLatest: hasLatest, Cover: cover, HasCover: hasCover}, nil
|
||||
}
|
||||
Reference in New Issue
Block a user