Co-authored-by: Sulthan Zaki <sultankiki05@gmail.com> Co-committed-by: Sulthan Zaki <sultankiki05@gmail.com>
This commit was merged in pull request #29.
This commit is contained in:
@@ -23,9 +23,9 @@ type Fetcher interface {
|
||||
// Two clocks, deliberately independent:
|
||||
//
|
||||
// - Interval is how often this goroutine wakes up and looks.
|
||||
// - Cooldown is how long one bookmark rests since its own last check.
|
||||
// - Cooldown is how long one series rests since its own last check.
|
||||
//
|
||||
// Only the cooldown is per bookmark, and it is enforced by the WHERE clause in
|
||||
// Only the cooldown is per series, and it is enforced by the WHERE clause in
|
||||
// DueForLatestCheck rather than by any timer. Shortening Interval therefore
|
||||
// cannot shorten anyone's cooldown; it only makes the poller wake up and find
|
||||
// nothing due more often.
|
||||
@@ -78,7 +78,7 @@ func (p *Poller) Run(ctx context.Context) {
|
||||
}
|
||||
}
|
||||
|
||||
// runOnce processes one batch of due bookmarks.
|
||||
// runOnce processes one batch of due series.
|
||||
func (p *Poller) runOnce(ctx context.Context) {
|
||||
cutoff := p.Now().Add(-p.Cooldown).UnixMilli()
|
||||
due, err := p.Store.DueForLatestCheck(cutoff, p.Batch)
|
||||
@@ -88,7 +88,7 @@ func (p *Poller) runOnce(ctx context.Context) {
|
||||
}
|
||||
|
||||
checked := 0
|
||||
for i, b := range due {
|
||||
for i, sr := range due {
|
||||
if ctx.Err() != nil {
|
||||
break
|
||||
}
|
||||
@@ -107,7 +107,7 @@ func (p *Poller) runOnce(ctx context.Context) {
|
||||
if stopped {
|
||||
break
|
||||
}
|
||||
p.checkOne(ctx, b)
|
||||
p.checkOne(ctx, sr)
|
||||
checked++
|
||||
}
|
||||
// due vs checked is how you tell which constraint is binding: ticks that
|
||||
@@ -119,10 +119,10 @@ func (p *Poller) runOnce(ctx context.Context) {
|
||||
// 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
|
||||
// batch or take down the process.
|
||||
func (p *Poller) checkOne(ctx context.Context, b store.Bookmark) {
|
||||
func (p *Poller) checkOne(ctx context.Context, sr store.Series) {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
log.Printf("latest poll %q: recovered from panic: %v", b.Key, r)
|
||||
log.Printf("latest poll %q: recovered from panic: %v", sr.Key(), r)
|
||||
}
|
||||
}()
|
||||
|
||||
@@ -130,8 +130,8 @@ func (p *Poller) checkOne(ctx context.Context, b store.Bookmark) {
|
||||
// mid-request still consumes the cooldown. Otherwise a renamed or deleted
|
||||
// series would be retried on every single tick forever. The userscript
|
||||
// stamps in the same order and for the same reason (L471-473).
|
||||
if err := p.Store.MarkLatestChecked(b.Key, p.Now().UnixMilli()); err != nil {
|
||||
log.Printf("latest poll %q: mark checked: %v", b.Key, err)
|
||||
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
|
||||
}
|
||||
|
||||
@@ -142,69 +142,52 @@ func (p *Poller) checkOne(ctx context.Context, b store.Bookmark) {
|
||||
// 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(b.Site, b.SeriesURL) {
|
||||
log.Printf("latest poll %q: not fetchable: site=%q url=%q", b.Key, b.Site, b.SeriesURL)
|
||||
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 := p.fetcherFor(b.Site)
|
||||
f := p.fetcherFor(sr.Site)
|
||||
if f == nil {
|
||||
log.Printf("latest poll %q: no fetcher for site %q", b.Key, b.Site)
|
||||
log.Printf("latest poll %q: no fetcher for site %q", sr.Key(), sr.Site)
|
||||
return
|
||||
}
|
||||
|
||||
body, status, err := f.Get(ctx, b.SeriesURL)
|
||||
body, status, err := f.Get(ctx, sr.SeriesURL)
|
||||
if err != nil {
|
||||
log.Printf("latest poll %q: fetch %s: %v", b.Key, b.SeriesURL, err)
|
||||
log.Printf("latest poll %q: fetch %s: %v", sr.Key(), sr.SeriesURL, err)
|
||||
return
|
||||
}
|
||||
if status != 200 {
|
||||
log.Printf("latest poll %q: fetch %s: status %d", b.Key, b.SeriesURL, status)
|
||||
log.Printf("latest poll %q: fetch %s: status %d", sr.Key(), sr.SeriesURL, status)
|
||||
return
|
||||
}
|
||||
|
||||
latest, ok := latestChapterFrom(b.Site, b.SeriesURL, body)
|
||||
latest, ok := latestChapterFrom(sr.Site, sr.SeriesURL, body)
|
||||
if !ok {
|
||||
// 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", b.Key, len(body))
|
||||
log.Printf("latest poll %q: no chapter links in %d bytes", sr.Key(), len(body))
|
||||
return
|
||||
}
|
||||
|
||||
// Re-read: the row may have been updated or deleted while the fetch was in
|
||||
// flight, and writing b back wholesale would undo that.
|
||||
//
|
||||
// ponytail: non-transactional read-modify-write, wrap Get+Upsert in a tx if
|
||||
// this ever runs for more than one user. A client PUT that commits between
|
||||
// these two statements is lost to the stale re-read — reverting read
|
||||
// progress or a status change, and moving updated_at because the stored
|
||||
// value now differs. Accepted for a single-user deployment: the window is
|
||||
// milliseconds and the loser is one poll cycle.
|
||||
cur, found, err := p.Store.Get(b.Key)
|
||||
if err != nil {
|
||||
log.Printf("latest poll %q: reread: %v", b.Key, err)
|
||||
return
|
||||
}
|
||||
if !found {
|
||||
return
|
||||
}
|
||||
// Equality, not >, mirroring the userscript (L427): a site that retracts a
|
||||
// chapter should correct the stored number downward.
|
||||
if cur.LatestChapterNum != nil && *cur.LatestChapterNum == latest.Num {
|
||||
// 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 {
|
||||
return
|
||||
}
|
||||
|
||||
num := latest.Num
|
||||
cur.LatestChapter = latest.Label
|
||||
cur.LatestChapterNum = &num
|
||||
// A candidate only. last_chapter_num is untouched, so the CASE in Upsert
|
||||
// keeps the stored updated_at and the bookmark list does not reorder.
|
||||
cur.UpdatedAt = p.Now().UnixMilli()
|
||||
if _, err := p.Store.Upsert(cur); err != nil {
|
||||
log.Printf("latest poll %q: upsert: %v", b.Key, err)
|
||||
// Series-level write: the row is shared, so one update refreshes every
|
||||
// bookmark joining to it, and the bookmark's updated_at is never touched —
|
||||
// a newly published chapter is not reading progress and must not reorder
|
||||
// the list.
|
||||
if err := p.Store.SetLatestChapter(sr.Site, sr.SeriesID, latest.Label, 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", b.Key, latest.Label)
|
||||
log.Printf("latest poll %q: latest is now %s", sr.Key(), latest.Label)
|
||||
}
|
||||
|
||||
// fetchableSeriesURL reports whether site is a site latestChapterFrom knows how
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -25,26 +26,35 @@ func newTestStore(t *testing.T) *store.Store {
|
||||
return s
|
||||
}
|
||||
|
||||
// seedForCheck inserts a bookmark and forces its latest_checked_at.
|
||||
// seedForCheck inserts a bookmark (and with it its series) and forces the
|
||||
// series' latest_checked_at.
|
||||
func seedForCheck(t *testing.T, s *store.Store, key, seriesURL string, checkedAt int64) {
|
||||
t.Helper()
|
||||
site, seriesID, ok := strings.Cut(key, ":")
|
||||
if !ok {
|
||||
t.Fatalf("key %q: no ':' separator", key)
|
||||
}
|
||||
if _, err := s.Upsert(store.Bookmark{
|
||||
Key: key,
|
||||
Site: "asura",
|
||||
SeriesID: key,
|
||||
Site: site,
|
||||
SeriesID: seriesID,
|
||||
SeriesURL: seriesURL,
|
||||
UpdatedAt: 1000,
|
||||
}); err != nil {
|
||||
t.Fatalf("seed %q: %v", key, err)
|
||||
}
|
||||
if err := s.MarkLatestChecked(key, checkedAt); err != nil {
|
||||
if err := s.MarkLatestChecked(site, seriesID, checkedAt); err != nil {
|
||||
t.Fatalf("seed mark %q: %v", key, err)
|
||||
}
|
||||
}
|
||||
|
||||
func readLatestCheckedAt(t *testing.T, s *store.Store, key string) int64 {
|
||||
t.Helper()
|
||||
ts, err := s.LatestCheckedAt(key)
|
||||
site, seriesID, ok := strings.Cut(key, ":")
|
||||
if !ok {
|
||||
t.Fatalf("key %q: no ':' separator", key)
|
||||
}
|
||||
ts, err := s.LatestCheckedAt(site, seriesID)
|
||||
if err != nil {
|
||||
t.Fatalf("LatestCheckedAt %q: %v", key, err)
|
||||
}
|
||||
@@ -217,6 +227,41 @@ func TestRunOnceRespectsBatchLimit(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// The point of the split (ADR-0003): a series referenced by several bookmarks
|
||||
// is fetched once per due cycle, not once per bookmark. Today the bookmark key
|
||||
// is <site>:<series_id>, so the second bookmark only exists once keys stop
|
||||
// being derived from the series identity (issue #22).
|
||||
func TestRunOnceFetchesSharedSeriesOnce(t *testing.T) {
|
||||
s := newTestStore(t)
|
||||
// The slug must match the fixture's own anchors: asura's parser scopes
|
||||
// chapter links to the stored slug.
|
||||
const slug = "chronicles-of-the-demon-faction-f886a8af"
|
||||
const url = "https://asurascans.com/comics/" + slug
|
||||
seedForCheck(t, s, "asura:"+slug, url, 0)
|
||||
if _, err := s.Upsert(store.Bookmark{
|
||||
Key: "asura:" + slug + ":2", Site: "asura", SeriesID: slug, UpdatedAt: 2000,
|
||||
}); err != nil {
|
||||
t.Fatalf("seed second reader: %v", err)
|
||||
}
|
||||
|
||||
f := &fakeFetcher{body: asuraSeriesFixture, status: 200}
|
||||
newTestPoller(t, s, f, time.UnixMilli(5_000_000)).runOnce(context.Background())
|
||||
|
||||
if got := f.callCount(); got != 1 {
|
||||
t.Fatalf("fetched shared series %d times, want 1", got)
|
||||
}
|
||||
// Both bookmarks join to the same updated series row.
|
||||
for _, key := range []string{"asura:" + slug, "asura:" + slug + ":2"} {
|
||||
b, ok, err := s.Get(key)
|
||||
if err != nil || !ok {
|
||||
t.Fatalf("Get %s: %v ok=%v", key, err, ok)
|
||||
}
|
||||
if b.LatestChapterNum == nil || *b.LatestChapterNum != 181 {
|
||||
t.Fatalf("%s LatestChapterNum = %v, want 181", key, b.LatestChapterNum)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// One unreachable series must not abandon the rest of the batch.
|
||||
func TestRunOnceOneBadSeriesDoesNotStallBatch(t *testing.T) {
|
||||
s := newTestStore(t)
|
||||
@@ -330,8 +375,8 @@ func TestCheckOneValidatesSeriesURLBeforeFetching(t *testing.T) {
|
||||
|
||||
now := time.UnixMilli(4_000_000)
|
||||
f := &fakeFetcher{body: asuraSeriesFixture, status: 200}
|
||||
newTestPoller(t, s, f, now).checkOne(context.Background(), store.Bookmark{
|
||||
Key: key, Site: tt.site, SeriesURL: tt.seriesURL,
|
||||
newTestPoller(t, s, f, now).checkOne(context.Background(), store.Series{
|
||||
Site: tt.site, SeriesID: "x", SeriesURL: tt.seriesURL,
|
||||
})
|
||||
|
||||
if got := f.callCount(); got != tt.wantCalls {
|
||||
|
||||
Reference in New Issue
Block a user