feat(store): persist poll lane state

This commit is contained in:
2026-08-21 17:17:50 +07:00
parent 030ffdc26e
commit fd1131d11d
5 changed files with 391 additions and 5 deletions
+223 -3
View File
@@ -90,6 +90,34 @@ type Series struct {
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>"),
// used by the poller's logs and by tests asserting on the due queue.
func (s Series) Key() string { return s.Site + ":" + s.SeriesID }
@@ -188,9 +216,9 @@ const (
//go:embed migrations/*.sql
var migrations embed.FS
// bookmarkColumns is the only value ever concatenated into query text. It is a
// compile-time constant; every request value is bound as a parameter. The
// series-owned fields are joined in from the series table, in scanBookmark
// These column lists are the only values ever concatenated into query text.
// They are compile-time constants; every request value is bound as a parameter.
// 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).
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,
@@ -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,
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
// 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
@@ -613,6 +645,18 @@ func scanSeries(scan func(...any) error) (Series, error) {
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.
func (s *Store) Close() error { return s.db.Close() }
@@ -934,6 +978,182 @@ func (s *Store) Delete(readerID int64, key string) error {
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
// latest-chapter check has aged past cutoffMs, ordered by how many bookmarks
// reference them (descending) then least-recently-checked first. One Site per