feat(store): persist poll lane state
This commit is contained in:
+1
-2
@@ -31,8 +31,7 @@ image, TLS terminated by the reverse proxy so the service listens plain `:8080`.
|
|||||||
store must have that `TestMain` or it has no database at all.
|
store must have that `TestMain` or it has no database at all.
|
||||||
|
|
||||||
### Reader-owned store — `internal/store`, `internal/token`
|
### Reader-owned store — `internal/store`, `internal/token`
|
||||||
|
- The Reader-owned tables are `readers`, `bookmarks`, `series`, and `sessions`; auxiliary `covers`, `poll_lanes`, and `poll_passes` are also defined in the migrations.
|
||||||
Four tables; shape is in the migrations, behaviour in `Store`'s methods.
|
|
||||||
|
|
||||||
- **Credentials are derived, never stored.** `token.Token(TOKEN_KEY, discord_id, epoch)`
|
- **Credentials are derived, never stored.** `token.Token(TOKEN_KEY, discord_id, epoch)`
|
||||||
is an HMAC; only its SHA-256 reaches `readers.token_sha256`. So install URLs
|
is an HMAC; only its SHA-256 reaches `readers.token_sha256`. So install URLs
|
||||||
|
|||||||
@@ -0,0 +1,6 @@
|
|||||||
|
-- One durable state row per Poll Lane. Zero means no pause or refusal is set.
|
||||||
|
CREATE TABLE poll_lanes (
|
||||||
|
site text NOT NULL PRIMARY KEY,
|
||||||
|
paused_until bigint NOT NULL DEFAULT 0,
|
||||||
|
refuse_until bigint NOT NULL DEFAULT 0
|
||||||
|
);
|
||||||
@@ -0,0 +1,16 @@
|
|||||||
|
-- Append-only Lane Pass log. Timestamps are unix milliseconds from the poller's clock.
|
||||||
|
CREATE TABLE poll_passes (
|
||||||
|
site text NOT NULL,
|
||||||
|
ran_at bigint NOT NULL,
|
||||||
|
skip text NOT NULL,
|
||||||
|
due integer NOT NULL,
|
||||||
|
checked integer NOT NULL,
|
||||||
|
gap_ms bigint NOT NULL,
|
||||||
|
clamped boolean NOT NULL,
|
||||||
|
refused integer NOT NULL,
|
||||||
|
unreachable integer NOT NULL,
|
||||||
|
no_chapter integer NOT NULL,
|
||||||
|
unfetchable integer NOT NULL,
|
||||||
|
errors integer NOT NULL,
|
||||||
|
PRIMARY KEY (site, ran_at)
|
||||||
|
);
|
||||||
@@ -90,6 +90,34 @@ type Series struct {
|
|||||||
readerCount int
|
readerCount int
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// LanePass is one Poll Lane's durable pass snapshot. Pause and refusal stamps
|
||||||
|
// are joined from poll_lanes on read; they are not pass facts.
|
||||||
|
type LanePass struct {
|
||||||
|
Site string
|
||||||
|
RanAt int64
|
||||||
|
Skip string
|
||||||
|
Due, Checked int
|
||||||
|
GapMS int64
|
||||||
|
Clamped bool
|
||||||
|
Refused, Unreachable, NoChapter int
|
||||||
|
Unfetchable, Errors int
|
||||||
|
PausedUntil, RefuseUntil int64
|
||||||
|
}
|
||||||
|
|
||||||
|
// SiteOutcomes is one Site's summed Lane Pass outcomes over a caller-supplied
|
||||||
|
// window.
|
||||||
|
type SiteOutcomes struct {
|
||||||
|
Site string
|
||||||
|
Refused, Unreachable, NoChapter int
|
||||||
|
Unfetchable, Errors int
|
||||||
|
}
|
||||||
|
|
||||||
|
// LanePause is one persisted Lane pause stamp.
|
||||||
|
type LanePause struct {
|
||||||
|
Site string
|
||||||
|
PausedUntil int64
|
||||||
|
}
|
||||||
|
|
||||||
// Key returns the canonical identity in bookmark-key form ("<site>:<series_id>"),
|
// Key returns the canonical identity in bookmark-key form ("<site>:<series_id>"),
|
||||||
// used by the poller's logs and by tests asserting on the due queue.
|
// used by the poller's logs and by tests asserting on the due queue.
|
||||||
func (s Series) Key() string { return s.Site + ":" + s.SeriesID }
|
func (s Series) Key() string { return s.Site + ":" + s.SeriesID }
|
||||||
@@ -188,9 +216,9 @@ const (
|
|||||||
//go:embed migrations/*.sql
|
//go:embed migrations/*.sql
|
||||||
var migrations embed.FS
|
var migrations embed.FS
|
||||||
|
|
||||||
// bookmarkColumns is the only value ever concatenated into query text. It is a
|
// These column lists are the only values ever concatenated into query text.
|
||||||
// compile-time constant; every request value is bound as a parameter. The
|
// They are compile-time constants; every request value is bound as a parameter.
|
||||||
// series-owned fields are joined in from the series table, in scanBookmark
|
// The series-owned fields are joined in from the series table, in scanBookmark
|
||||||
// order, so the flat Bookmark reads back whole despite the split (ADR-0004).
|
// order, so the flat Bookmark reads back whole despite the split (ADR-0004).
|
||||||
const bookmarkColumns = `b.site, b.series_id, s.title, s.series_url, s.cover_address,
|
const bookmarkColumns = `b.site, b.series_id, s.title, s.series_url, s.cover_address,
|
||||||
b.last_chapter, b.last_chapter_num, b.last_chapter_url,
|
b.last_chapter, b.last_chapter_num, b.last_chapter_url,
|
||||||
@@ -202,6 +230,10 @@ const bookmarkColumns = `b.site, b.series_id, s.title, s.series_url, s.cover_add
|
|||||||
const seriesColumns = `s.site, s.series_id, s.title, s.series_url, s.cover, s.cover_address,
|
const seriesColumns = `s.site, s.series_id, s.title, s.series_url, s.cover, s.cover_address,
|
||||||
s.kind, s.latest_chapter, s.latest_chapter_num, s.latest_checked_at, s.latest_raised_by`
|
s.kind, s.latest_chapter, s.latest_chapter_num, s.latest_checked_at, s.latest_raised_by`
|
||||||
|
|
||||||
|
const lanePassColumns = `p.site, p.ran_at, p.skip, p.due, p.checked, p.gap_ms, p.clamped,
|
||||||
|
p.refused, p.unreachable, p.no_chapter, p.unfetchable, p.errors,
|
||||||
|
COALESCE(l.paused_until, 0), COALESCE(l.refuse_until, 0)`
|
||||||
|
|
||||||
// Owner is the person running the service: the first Reader, seeded at startup
|
// Owner is the person running the service: the first Reader, seeded at startup
|
||||||
// so a fresh deployment has a library before anyone logs in. The seed makes
|
// so a fresh deployment has a library before anyone logs in. The seed makes
|
||||||
// sure exactly one readers row matches their Discord ID, carrying the SHA-256
|
// sure exactly one readers row matches their Discord ID, carrying the SHA-256
|
||||||
@@ -613,6 +645,18 @@ func scanSeries(scan func(...any) error) (Series, error) {
|
|||||||
return sr, nil
|
return sr, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func scanLanePass(scan func(...any) error) (LanePass, error) {
|
||||||
|
var p LanePass
|
||||||
|
if err := scan(
|
||||||
|
&p.Site, &p.RanAt, &p.Skip, &p.Due, &p.Checked, &p.GapMS, &p.Clamped,
|
||||||
|
&p.Refused, &p.Unreachable, &p.NoChapter, &p.Unfetchable, &p.Errors,
|
||||||
|
&p.PausedUntil, &p.RefuseUntil,
|
||||||
|
); err != nil {
|
||||||
|
return LanePass{}, err
|
||||||
|
}
|
||||||
|
return p, nil
|
||||||
|
}
|
||||||
|
|
||||||
// Close releases the underlying database handle.
|
// Close releases the underlying database handle.
|
||||||
func (s *Store) Close() error { return s.db.Close() }
|
func (s *Store) Close() error { return s.db.Close() }
|
||||||
|
|
||||||
@@ -934,6 +978,182 @@ func (s *Store) Delete(readerID int64, key string) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// RecordLanePass appends one pass and prunes every older row in the same
|
||||||
|
// transaction. retainBefore is supplied by the poller's clock.
|
||||||
|
func (s *Store) RecordLanePass(p LanePass, retainBefore int64) error {
|
||||||
|
tx, err := s.db.Begin()
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("begin lane pass %s: %w", p.Site, err)
|
||||||
|
}
|
||||||
|
defer tx.Rollback()
|
||||||
|
if _, err := tx.Exec(`
|
||||||
|
INSERT INTO poll_passes
|
||||||
|
(site, ran_at, skip, due, checked, gap_ms, clamped,
|
||||||
|
refused, unreachable, no_chapter, unfetchable, errors)
|
||||||
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12)`,
|
||||||
|
p.Site, p.RanAt, p.Skip, p.Due, p.Checked, p.GapMS, p.Clamped,
|
||||||
|
p.Refused, p.Unreachable, p.NoChapter, p.Unfetchable, p.Errors); err != nil {
|
||||||
|
return fmt.Errorf("insert lane pass %s at %d: %w", p.Site, p.RanAt, err)
|
||||||
|
}
|
||||||
|
if _, err := tx.Exec(`DELETE FROM poll_passes WHERE ran_at < $1`, retainBefore); err != nil {
|
||||||
|
return fmt.Errorf("prune lane passes before %d: %w", retainBefore, err)
|
||||||
|
}
|
||||||
|
if err := tx.Commit(); err != nil {
|
||||||
|
return fmt.Errorf("commit lane pass %s at %d: %w", p.Site, p.RanAt, err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// LatestLanePass returns the newest pass for one Site, with its current Lane
|
||||||
|
// state joined on. A Site without a pass has no durable snapshot yet.
|
||||||
|
func (s *Store) LatestLanePass(site string) (LanePass, bool, error) {
|
||||||
|
p, err := scanLanePass(s.db.QueryRow(`SELECT `+lanePassColumns+`
|
||||||
|
FROM poll_passes p
|
||||||
|
LEFT JOIN poll_lanes l ON l.site = p.site
|
||||||
|
WHERE p.site = $1
|
||||||
|
ORDER BY p.ran_at DESC
|
||||||
|
LIMIT 1`, site).Scan)
|
||||||
|
if errors.Is(err, sql.ErrNoRows) {
|
||||||
|
return LanePass{}, false, nil
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
return LanePass{}, false, fmt.Errorf("latest lane pass %s: %w", site, err)
|
||||||
|
}
|
||||||
|
return p, true, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// LatestLanePasses returns the newest pass for each Site, with current Lane
|
||||||
|
// state joined on. Sites without a pass have no row yet.
|
||||||
|
func (s *Store) LatestLanePasses() ([]LanePass, error) {
|
||||||
|
rows, err := s.db.Query(`SELECT ` + lanePassColumns + `
|
||||||
|
FROM (
|
||||||
|
SELECT DISTINCT ON (site)
|
||||||
|
site, ran_at, skip, due, checked, gap_ms, clamped,
|
||||||
|
refused, unreachable, no_chapter, unfetchable, errors
|
||||||
|
FROM poll_passes
|
||||||
|
ORDER BY site, ran_at DESC
|
||||||
|
) p
|
||||||
|
LEFT JOIN poll_lanes l ON l.site = p.site
|
||||||
|
ORDER BY p.site`)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("query latest lane passes: %w", err)
|
||||||
|
}
|
||||||
|
defer rows.Close()
|
||||||
|
|
||||||
|
out := []LanePass{}
|
||||||
|
for rows.Next() {
|
||||||
|
p, err := scanLanePass(rows.Scan)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("scan latest lane pass: %w", err)
|
||||||
|
}
|
||||||
|
out = append(out, p)
|
||||||
|
}
|
||||||
|
return out, rows.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
// LanePassOutcomes sums the named outcomes for each Site at or after since.
|
||||||
|
// The window boundary is supplied by the caller; the store has no clock.
|
||||||
|
func (s *Store) LanePassOutcomes(since int64) ([]SiteOutcomes, error) {
|
||||||
|
rows, err := s.db.Query(`
|
||||||
|
SELECT site, SUM(refused), SUM(unreachable), SUM(no_chapter),
|
||||||
|
SUM(unfetchable), SUM(errors)
|
||||||
|
FROM poll_passes
|
||||||
|
WHERE ran_at >= $1
|
||||||
|
GROUP BY site
|
||||||
|
ORDER BY site`, since)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("query lane pass outcomes: %w", err)
|
||||||
|
}
|
||||||
|
defer rows.Close()
|
||||||
|
|
||||||
|
out := []SiteOutcomes{}
|
||||||
|
for rows.Next() {
|
||||||
|
var outcomes SiteOutcomes
|
||||||
|
if err := rows.Scan(
|
||||||
|
&outcomes.Site, &outcomes.Refused, &outcomes.Unreachable,
|
||||||
|
&outcomes.NoChapter, &outcomes.Unfetchable, &outcomes.Errors,
|
||||||
|
); err != nil {
|
||||||
|
return nil, fmt.Errorf("scan lane pass outcomes: %w", err)
|
||||||
|
}
|
||||||
|
out = append(out, outcomes)
|
||||||
|
}
|
||||||
|
return out, rows.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
// SetLaneRefusal persists a Site's refusal backoff stamp without touching its
|
||||||
|
// pause. until is supplied by the caller's clock.
|
||||||
|
func (s *Store) SetLaneRefusal(site string, until int64) error {
|
||||||
|
if _, err := s.db.Exec(`
|
||||||
|
INSERT INTO poll_lanes (site, refuse_until) VALUES ($1, $2)
|
||||||
|
ON CONFLICT (site) DO UPDATE SET refuse_until = EXCLUDED.refuse_until`, site, until); err != nil {
|
||||||
|
return fmt.Errorf("set lane refusal %s: %w", site, err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// PauseLane persists a bounded pause. The caller must ensure until is after
|
||||||
|
// its current timestamp; the store has no clock and rejects only the invalid
|
||||||
|
// zero and negative sentinels.
|
||||||
|
func (s *Store) PauseLane(site string, until int64) error {
|
||||||
|
if until <= 0 {
|
||||||
|
return fmt.Errorf("pause lane %s: expiry must be positive", site)
|
||||||
|
}
|
||||||
|
if _, err := s.db.Exec(`
|
||||||
|
INSERT INTO poll_lanes (site, paused_until) VALUES ($1, $2)
|
||||||
|
ON CONFLICT (site) DO UPDATE SET paused_until = EXCLUDED.paused_until`, site, until); err != nil {
|
||||||
|
return fmt.Errorf("pause lane %s: %w", site, err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ResumeLane clears only the pause stamp and keeps the Lane state row, along
|
||||||
|
// with any refusal stamp already persisted on it.
|
||||||
|
func (s *Store) ResumeLane(site string) error {
|
||||||
|
if _, err := s.db.Exec(
|
||||||
|
`UPDATE poll_lanes SET paused_until = 0 WHERE site = $1`, site); err != nil {
|
||||||
|
return fmt.Errorf("resume lane %s: %w", site, err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// PausedLanes returns Lane rows with a nonzero pause stamp. Expiry comparison
|
||||||
|
// stays with the caller because the store is deliberately clockless.
|
||||||
|
func (s *Store) PausedLanes() ([]LanePause, error) {
|
||||||
|
rows, err := s.db.Query(`
|
||||||
|
SELECT site, paused_until
|
||||||
|
FROM poll_lanes
|
||||||
|
WHERE paused_until > 0
|
||||||
|
ORDER BY site`)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("query paused lanes: %w", err)
|
||||||
|
}
|
||||||
|
defer rows.Close()
|
||||||
|
|
||||||
|
out := []LanePause{}
|
||||||
|
for rows.Next() {
|
||||||
|
var pause LanePause
|
||||||
|
if err := rows.Scan(&pause.Site, &pause.PausedUntil); err != nil {
|
||||||
|
return nil, fmt.Errorf("scan paused lane: %w", err)
|
||||||
|
}
|
||||||
|
out = append(out, pause)
|
||||||
|
}
|
||||||
|
return out, rows.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
// LanePausedUntil reads a Site's pause stamp. A missing state row is the
|
||||||
|
// default unpaused state.
|
||||||
|
func (s *Store) LanePausedUntil(site string) (int64, error) {
|
||||||
|
var until int64
|
||||||
|
err := s.db.QueryRow(`SELECT paused_until FROM poll_lanes WHERE site = $1`, site).Scan(&until)
|
||||||
|
if errors.Is(err, sql.ErrNoRows) {
|
||||||
|
return 0, nil
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("lane pause %s: %w", site, err)
|
||||||
|
}
|
||||||
|
return until, nil
|
||||||
|
}
|
||||||
|
|
||||||
// DueForLatestCheck returns one Site's series whose server-side
|
// DueForLatestCheck returns one Site's series whose server-side
|
||||||
// latest-chapter check has aged past cutoffMs, ordered by how many bookmarks
|
// latest-chapter check has aged past cutoffMs, ordered by how many bookmarks
|
||||||
// reference them (descending) then least-recently-checked first. One Site per
|
// reference them (descending) then least-recently-checked first. One Site per
|
||||||
|
|||||||
@@ -1584,3 +1584,148 @@ func TestCoverStoreAcceptsAnySourceURL(t *testing.T) {
|
|||||||
t.Fatalf("rejected cover = found %v, err %v; want missing", ok, err)
|
t.Fatalf("rejected cover = found %v, err %v; want missing", ok, err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
func TestRecordLanePassPrunesBeforeInsertCutoff(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
for _, pass := range []LanePass{
|
||||||
|
{Site: "asura", RanAt: 99},
|
||||||
|
{Site: "asura", RanAt: 100},
|
||||||
|
} {
|
||||||
|
if err := s.RecordLanePass(pass, 100); err != nil {
|
||||||
|
t.Fatalf("RecordLanePass(%d): %v", pass.RanAt, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := s.RecordLanePass(LanePass{Site: "asura", RanAt: 200}, 100); err != nil {
|
||||||
|
t.Fatalf("RecordLanePass(200): %v", err)
|
||||||
|
}
|
||||||
|
var count int
|
||||||
|
if err := s.db.QueryRow(`SELECT count(*) FROM poll_passes WHERE site = $1`, "asura").Scan(&count); err != nil {
|
||||||
|
t.Fatalf("count passes: %v", err)
|
||||||
|
}
|
||||||
|
if count != 2 {
|
||||||
|
t.Fatalf("retained passes = %d, want 2", count)
|
||||||
|
}
|
||||||
|
if _, ok, err := s.LatestLanePass("asura"); err != nil || !ok {
|
||||||
|
t.Fatalf("LatestLanePass = ok %v, err %v; want latest row", ok, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestLatestLanePassesKeepsNewestPerSiteAndJoinsState(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
for _, pass := range []LanePass{
|
||||||
|
{Site: "asura", RanAt: 100, Due: 1},
|
||||||
|
{Site: "asura", RanAt: 200, Skip: "due-query", Due: 2, Checked: 3, GapMS: 4000, Clamped: true},
|
||||||
|
{Site: "demonic", RanAt: 150, Due: 4},
|
||||||
|
} {
|
||||||
|
if err := s.RecordLanePass(pass, -1); err != nil {
|
||||||
|
t.Fatalf("RecordLanePass(%s/%d): %v", pass.Site, pass.RanAt, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if err := s.PauseLane("asura", 1234); err != nil {
|
||||||
|
t.Fatalf("PauseLane: %v", err)
|
||||||
|
}
|
||||||
|
if err := s.SetLaneRefusal("asura", 5678); err != nil {
|
||||||
|
t.Fatalf("SetLaneRefusal: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
got, err := s.LatestLanePasses()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("LatestLanePasses: %v", err)
|
||||||
|
}
|
||||||
|
if len(got) != 2 {
|
||||||
|
t.Fatalf("latest passes = %d, want one per Site", len(got))
|
||||||
|
}
|
||||||
|
bySite := map[string]LanePass{}
|
||||||
|
for _, pass := range got {
|
||||||
|
bySite[pass.Site] = pass
|
||||||
|
}
|
||||||
|
asura := bySite["asura"]
|
||||||
|
if asura.RanAt != 200 || asura.Skip != "due-query" || asura.Due != 2 || asura.Checked != 3 || asura.GapMS != 4000 || !asura.Clamped ||
|
||||||
|
asura.PausedUntil != 1234 || asura.RefuseUntil != 5678 {
|
||||||
|
t.Fatalf("asura latest pass = %+v, want newest pass and joined state", asura)
|
||||||
|
}
|
||||||
|
if demonic := bySite["demonic"]; demonic.RanAt != 150 || demonic.Due != 4 {
|
||||||
|
t.Fatalf("demonic latest pass = %+v, want its only pass", demonic)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestLanePassOutcomesSumsWindow(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
for _, pass := range []LanePass{
|
||||||
|
{Site: "asura", RanAt: 99, Refused: 1, Unreachable: 2, NoChapter: 3, Unfetchable: 4, Errors: 5},
|
||||||
|
{Site: "asura", RanAt: 100, Refused: 2, Unreachable: 3, NoChapter: 4, Unfetchable: 5, Errors: 6},
|
||||||
|
{Site: "asura", RanAt: 200, Refused: 3, Unreachable: 4, NoChapter: 5, Unfetchable: 6, Errors: 7},
|
||||||
|
{Site: "demonic", RanAt: 150, Refused: 8, Unreachable: 9, NoChapter: 10, Unfetchable: 11, Errors: 12},
|
||||||
|
} {
|
||||||
|
if err := s.RecordLanePass(pass, -1); err != nil {
|
||||||
|
t.Fatalf("RecordLanePass(%s/%d): %v", pass.Site, pass.RanAt, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
got, err := s.LanePassOutcomes(100)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("LanePassOutcomes: %v", err)
|
||||||
|
}
|
||||||
|
if len(got) != 2 {
|
||||||
|
t.Fatalf("outcome Sites = %d, want 2", len(got))
|
||||||
|
}
|
||||||
|
bySite := map[string]SiteOutcomes{}
|
||||||
|
for _, outcomes := range got {
|
||||||
|
bySite[outcomes.Site] = outcomes
|
||||||
|
}
|
||||||
|
if want := (SiteOutcomes{Site: "asura", Refused: 5, Unreachable: 7, NoChapter: 9, Unfetchable: 11, Errors: 13}); bySite["asura"] != want {
|
||||||
|
t.Fatalf("asura outcomes = %+v, want %+v", bySite["asura"], want)
|
||||||
|
}
|
||||||
|
if want := (SiteOutcomes{Site: "demonic", Refused: 8, Unreachable: 9, NoChapter: 10, Unfetchable: 11, Errors: 12}); bySite["demonic"] != want {
|
||||||
|
t.Fatalf("demonic outcomes = %+v, want %+v", bySite["demonic"], want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestLaneStatePauseResumeAndRefusal(t *testing.T) {
|
||||||
|
s := newTestStore(t)
|
||||||
|
for _, until := range []int64{0, -1} {
|
||||||
|
if err := s.PauseLane("asura", until); err == nil {
|
||||||
|
t.Fatalf("PauseLane(%d) accepted a non-future expiry", until)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if err := s.PauseLane("asura", 2000); err != nil {
|
||||||
|
t.Fatalf("PauseLane: %v", err)
|
||||||
|
}
|
||||||
|
if err := s.SetLaneRefusal("asura", 3000); err != nil {
|
||||||
|
t.Fatalf("SetLaneRefusal: %v", err)
|
||||||
|
}
|
||||||
|
if err := s.RecordLanePass(LanePass{Site: "asura", RanAt: 1}, -1); err != nil {
|
||||||
|
t.Fatalf("RecordLanePass: %v", err)
|
||||||
|
}
|
||||||
|
if got, err := s.LanePausedUntil("asura"); err != nil || got != 2000 {
|
||||||
|
t.Fatalf("LanePausedUntil = %d, %v; want 2000", got, err)
|
||||||
|
}
|
||||||
|
paused, err := s.PausedLanes()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("PausedLanes: %v", err)
|
||||||
|
}
|
||||||
|
if len(paused) != 1 || paused[0] != (LanePause{Site: "asura", PausedUntil: 2000}) {
|
||||||
|
t.Fatalf("PausedLanes = %+v, want asura/2000", paused)
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := s.ResumeLane("asura"); err != nil {
|
||||||
|
t.Fatalf("ResumeLane: %v", err)
|
||||||
|
}
|
||||||
|
latest, ok, err := s.LatestLanePass("asura")
|
||||||
|
if err != nil || !ok || latest.RefuseUntil != 3000 {
|
||||||
|
t.Fatalf("latest refusal after resume = %+v, ok=%v, err=%v; want 3000 preserved", latest, ok, err)
|
||||||
|
}
|
||||||
|
if got, err := s.LanePausedUntil("asura"); err != nil || got != 0 {
|
||||||
|
t.Fatalf("LanePausedUntil after resume = %d, %v; want 0", got, err)
|
||||||
|
}
|
||||||
|
if paused, err := s.PausedLanes(); err != nil || len(paused) != 0 {
|
||||||
|
t.Fatalf("PausedLanes after resume = %+v, %v; want empty", paused, err)
|
||||||
|
}
|
||||||
|
var rows int
|
||||||
|
if err := s.db.QueryRow(`SELECT count(*) FROM poll_lanes WHERE site = $1`, "asura").Scan(&rows); err != nil {
|
||||||
|
t.Fatalf("count lane state: %v", err)
|
||||||
|
}
|
||||||
|
if rows != 1 {
|
||||||
|
t.Fatalf("lane state rows after resume = %d, want 1", rows)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user