bc64a1d894
(*Store).RemoveSeries deletes one series row by (site, series_id) via a
plain parameterized DELETE; a bookmarks_series_fk violation outside
23503 is translated into the ErrSeriesHasBookmarks sentinel so no
driver type escapes the store. The caller reads the row's cover
address before the delete and reclaims it after: ReclaimCover's guard
cannot pass while a series row still points at the address.
POST /admin/series/{key}/remove answers the list row with the removed
row's fragment plus the heading re-rendered out of band at the fresh
count (HX-Reswap: delete removes the row through the same button that
swaps the refusal back in), and navigates from the detail page to the
No-Readers list (HX-Redirect for htmx, a 303 for plain clients). A
removal that races a fresh Bookmark is a refusal, not a 500: the row
re-renders at its new count with the fact spelled out. The control
renders only at zero Reader count on both surfaces, gated by hx-confirm
with the brief's copy.
1608 lines
65 KiB
Go
1608 lines
65 KiB
Go
package store
|
|
|
|
import (
|
|
"crypto/sha256"
|
|
"database/sql"
|
|
"embed"
|
|
"encoding/hex"
|
|
"errors"
|
|
"fmt"
|
|
"io/fs"
|
|
"os"
|
|
"path"
|
|
"path/filepath"
|
|
"regexp"
|
|
"slices"
|
|
"strconv"
|
|
"strings"
|
|
|
|
"github.com/jackc/pgx/v5/pgconn"
|
|
_ "github.com/jackc/pgx/v5/stdlib"
|
|
)
|
|
|
|
// Bookmark is one tracked series, keyed "<site>:<series_id>" across both sites.
|
|
//
|
|
// LastChapter* is the user's read progress; LatestChapter* is the newest
|
|
// chapter the site has published, captured opportunistically by the userscript.
|
|
//
|
|
// Title, SeriesURL, Cover, Kind and LatestChapter* live on the shared Series
|
|
// row (ADR-0003) and are joined in on read; Bookmark carries only what differs
|
|
// between readers: progress, favourite, lifecycle bucket, updated_at. The wire
|
|
// format stays flat regardless — see ADR-0004.
|
|
type Bookmark struct {
|
|
Key string `json:"key"`
|
|
Site string `json:"site"`
|
|
SeriesID string `json:"series_id"`
|
|
Title string `json:"title"`
|
|
SeriesURL string `json:"series_url"`
|
|
// Cover is the wire value: an absolute URL on this deployment's own
|
|
// origin once the bytes exist, and "" until they do — never a third-party
|
|
// address and never an address that 404s (ADR-0007). A client may still
|
|
// send this field and it is discarded on the way in; see Upsert.
|
|
Cover string `json:"cover"`
|
|
LastChapter string `json:"last_chapter"`
|
|
LastChapterNum float64 `json:"last_chapter_num"`
|
|
LastChapterURL string `json:"last_chapter_url"`
|
|
Favorite bool `json:"favorite"`
|
|
LatestChapter string `json:"latest_chapter"`
|
|
LatestChapterNum *float64 `json:"latest_chapter_num"` // nil until first captured
|
|
UpdatedAt int64 `json:"updated_at"` // unix ms; see Upsert
|
|
// Status is the lifecycle bucket: reading, archived, or finished.
|
|
// Archived series stay polled for new chapters; finished ones do not.
|
|
// Empty on the way in means "no opinion" — see Upsert.
|
|
Status string `json:"status"`
|
|
// Kind is the library bucket: manga or novel. Empty on the way in means
|
|
// "no opinion" — see Upsert.
|
|
Kind string `json:"kind"`
|
|
}
|
|
|
|
// Series is one distinct work, shared by every bookmark that tracks it. It is
|
|
// keyed (site, series_id) — the pair a bookmark key decomposes into — and
|
|
// exists once no matter how many bookmarks point at it (ADR-0003).
|
|
//
|
|
// Title, SeriesURL and Cover are written once, at creation: a PUT naming an
|
|
// existing Series has them ignored, and only the backend's own Poll may
|
|
// change them. The one exception is SeriesURL, which the owner's
|
|
// SetSeriesURL may repair (issue #151). Kind and the latest-chapter fields
|
|
// are last-write-wins like the bookmark's own fields. Never serialized: the
|
|
// wire format is the flat Bookmark (ADR-0004).
|
|
type Series struct {
|
|
Site string
|
|
SeriesID string
|
|
Title string
|
|
SeriesURL string
|
|
// Cover is the third-party source address the bytes come from, and
|
|
// CoverAddress the content address they are stored under. A blank
|
|
// CoverAddress is what "no Cover yet" means: the poll fills it and never
|
|
// replaces a filled one (ADR-0007).
|
|
Cover string
|
|
CoverAddress string
|
|
Kind string
|
|
LatestChapter string
|
|
LatestChapterNum *float64 // nil until first captured
|
|
LatestCheckedAt int64 // unix ms; see MarkLatestChecked
|
|
// LatestRaisedBy is the Reader whose Sighting last raised LatestChapter,
|
|
// and nil when the stored value is a Poll's own finding. It is what lets a
|
|
// Poll that contradicts the value downwards name a Reader instead of
|
|
// merely flagging the row (issue #103); the Poll that judges it clears it.
|
|
LatestRaisedBy *int64
|
|
|
|
// readerCount is the number of bookmarks referencing this series, filled
|
|
// only by the due-queue query that orders on it.
|
|
readerCount int
|
|
// Forced is whether the owner asked for a check now (issue #146): the
|
|
// request stamp is newer than the check stamp. Derived in the due query,
|
|
// never stored, and the flag that jumps the queue and opens the browser
|
|
// wake gate.
|
|
Forced bool
|
|
}
|
|
|
|
// 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 }
|
|
|
|
// HasNewChapter reports whether the site has published past the read point.
|
|
// A nil LatestChapterNum means nothing has been captured yet, which is not the
|
|
// same as "nothing new".
|
|
func (b Bookmark) HasNewChapter() bool {
|
|
return b.LatestChapterNum != nil && *b.LatestChapterNum > b.LastChapterNum
|
|
}
|
|
|
|
// chapterLeadIn matches the prefix the userscript and the poller both write
|
|
// ("Chapter 250"), so the UI can add exactly one "Ch " of its own instead of
|
|
// doubling it. A manual edit through the web UI stores a bare "250", which is
|
|
// the same string minus the lead-in.
|
|
var chapterLeadIn = regexp.MustCompile(`(?i)^\s*(?:chapter|ch\.?)\s*`)
|
|
|
|
func displayChapter(raw string, num float64) string {
|
|
rest := strings.TrimSpace(chapterLeadIn.ReplaceAllString(raw, ""))
|
|
if rest == "" {
|
|
rest = strconv.FormatFloat(num, 'f', -1, 64)
|
|
}
|
|
// "Ch " only makes sense in front of a number; anything else is a label the
|
|
// site gave us, so pass it through as written.
|
|
if rest[0] < '0' || rest[0] > '9' {
|
|
return rest
|
|
}
|
|
return "Ch " + rest
|
|
}
|
|
|
|
// DisplayChapter is the read-progress line: one canonical "Ch N" whatever
|
|
// format the write came in as.
|
|
func (b Bookmark) DisplayChapter() string {
|
|
return displayChapter(b.LastChapter, b.LastChapterNum)
|
|
}
|
|
|
|
// DisplayLatest is the same for the newest published chapter, which arrives
|
|
// with the same "Chapter N" lead-in from both the userscript and the poller.
|
|
func (b Bookmark) DisplayLatest() string {
|
|
var num float64
|
|
if b.LatestChapterNum != nil {
|
|
num = *b.LatestChapterNum
|
|
}
|
|
return displayChapter(b.LatestChapter, num)
|
|
}
|
|
|
|
// ContinueURL is where the Continue button points: the chapter last read, or
|
|
// the series page when no chapter URL was ever captured.
|
|
func (b Bookmark) ContinueURL() string {
|
|
if b.LastChapterURL != "" {
|
|
return b.LastChapterURL
|
|
}
|
|
return b.SeriesURL
|
|
}
|
|
|
|
// Initial is the monogram the web UI shows in place of a cover when the
|
|
// source site never gave us an og:image. First rune, uppercased; "?" when even
|
|
// the title is missing, so the slot is never empty.
|
|
func (b Bookmark) Initial() string {
|
|
for _, r := range b.Title {
|
|
return strings.ToUpper(string(r))
|
|
}
|
|
return "?"
|
|
}
|
|
|
|
// CoverContentType canonicalises a fetched response's media type and reports
|
|
// whether the bytes are safe to store and serve. comix answers "image/jpg",
|
|
// which no standard lists but browsers accept; it is stored as the real name
|
|
// rather than passed through, so one image never lands under two spellings.
|
|
func CoverContentType(contentType string) (string, bool) {
|
|
switch contentType {
|
|
case "image/jpg":
|
|
return "image/jpeg", true
|
|
case "image/webp", "image/jpeg", "image/png", "image/avif", "image/gif":
|
|
return contentType, true
|
|
default:
|
|
return "", false
|
|
}
|
|
}
|
|
|
|
// Library buckets. A bookmark is in exactly one. This cannot be derived from
|
|
// Site: asurascans serves manga and novels from the same /comics/ path, so the
|
|
// userscript that recorded the page is the only party that knows which.
|
|
const (
|
|
KindManga = "manga"
|
|
KindNovel = "novel"
|
|
)
|
|
|
|
// Lifecycle buckets. A bookmark is in exactly one; favorite is orthogonal.
|
|
const (
|
|
StatusReading = "reading"
|
|
StatusArchived = "archived"
|
|
StatusFinished = "finished"
|
|
)
|
|
|
|
//go:embed migrations/*.sql
|
|
var migrations embed.FS
|
|
|
|
// 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,
|
|
b.favorite, s.latest_chapter, s.latest_chapter_num, b.updated_at, b.status, s.kind`
|
|
|
|
// seriesColumns is the series row in scanSeries order, used by the poller's
|
|
// due query. latest_checked_at lives only on series — see MarkLatestChecked
|
|
// for why it stays off every client-visible write.
|
|
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
|
|
// of their epoch-0 userscript credential (derived by internal/token). Every
|
|
// other Reader is created by their own first login (EnsureReader).
|
|
type Owner struct {
|
|
DiscordID string
|
|
// TokenHash is the SHA-256 of the epoch-0 credential; the array shape
|
|
// makes it a compile error to store anything that is not a hash.
|
|
TokenHash [32]byte
|
|
}
|
|
|
|
// Store is the Postgres-backed bookmark store.
|
|
type Store struct {
|
|
db *sql.DB
|
|
// ownerID is the seeded owner Reader (issue #22) — the only Reader with
|
|
// administrative reach (revoking another Reader's sessions). Every store
|
|
// method takes a reader id explicitly, so ownership is never implicit.
|
|
ownerID int64
|
|
coverDir string
|
|
// coverBaseURL is this deployment's public origin. Cover addresses are
|
|
// absolute because the userscript renders them on third-party origins,
|
|
// where a relative path would resolve against the Site (ADR-0007).
|
|
coverBaseURL string
|
|
// OnSeriesCreated fires once, after commit, for a Series no Reader had
|
|
// bookmarked before. It is how creation-time Cover and Latest Chapter
|
|
// acquisition is triggered without the write waiting on a third-party
|
|
// Site; nil disables it, which is what every test that does not care
|
|
// about acquisition leaves it as.
|
|
OnSeriesCreated func(Series)
|
|
}
|
|
|
|
// OwnerID returns the seeded owner Reader's id: the administrator, and the
|
|
// Reader every pre-registration bookmark belongs to.
|
|
func (s *Store) OwnerID() int64 { return s.ownerID }
|
|
|
|
// ReaderIDForTokenHash resolves the Reader whose stored credential hash
|
|
// matches, reporting absence with ok=false. The comparison is an equality on
|
|
// the 32-byte SHA-256 of the presented credential — never on the credential
|
|
// itself — and the indexed lookup reveals only whether some Reader matches,
|
|
// which the 401/200 split has to reveal anyway. An attacker's probe is the
|
|
// hash of their guess, so even the index's prefix comparisons leak nothing
|
|
// about the real credential.
|
|
func (s *Store) ReaderIDForTokenHash(hash [32]byte) (int64, bool, error) {
|
|
var id int64
|
|
err := s.db.QueryRow(
|
|
`SELECT id FROM readers WHERE token_sha256 = $1`, hash[:]).Scan(&id)
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return 0, false, nil
|
|
}
|
|
if err != nil {
|
|
return 0, false, fmt.Errorf("reader by token hash: %w", err)
|
|
}
|
|
return id, true, nil
|
|
}
|
|
|
|
// ReaderTokenInfo returns the identity halves a Reader's credential is
|
|
// derived from (internal/token.Token): their Discord id and token epoch. The
|
|
// web UI needs these to rebuild the install URL — the only place a credential
|
|
// is ever produced in plaintext.
|
|
func (s *Store) ReaderTokenInfo(readerID int64) (string, int64, error) {
|
|
var (
|
|
discordID string
|
|
epoch int64
|
|
)
|
|
err := s.db.QueryRow(
|
|
`SELECT discord_id, token_epoch FROM readers WHERE id = $1`, readerID).
|
|
Scan(&discordID, &epoch)
|
|
if err != nil {
|
|
return "", 0, fmt.Errorf("reader %d token info: %w", readerID, err)
|
|
}
|
|
return discordID, epoch, nil
|
|
}
|
|
|
|
// RotateToken bumps a Reader's token epoch and rewrites the stored hash in
|
|
// one statement, so the new hash always matches the new epoch. expectedEpoch
|
|
// is the epoch the caller derived newHash for (ReaderTokenInfo + 1); a
|
|
// concurrent rotation — or an unknown reader — leaves the row untouched and
|
|
// is reported as an error rather than silently succeeding.
|
|
func (s *Store) RotateToken(readerID, expectedEpoch int64, newHash [32]byte) error {
|
|
var epoch int64
|
|
err := s.db.QueryRow(`
|
|
UPDATE readers SET token_epoch = token_epoch + 1, token_sha256 = $3
|
|
WHERE id = $1 AND token_epoch = $2
|
|
RETURNING token_epoch`, readerID, expectedEpoch, newHash[:]).Scan(&epoch)
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return fmt.Errorf("rotate token for reader %d: concurrent rotation or unknown reader", readerID)
|
|
}
|
|
if err != nil {
|
|
return fmt.Errorf("rotate token for reader %d: %w", readerID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// EnsureReader returns the Reader registered to discordID, creating the row on
|
|
// first sight. Registration is open to every guild member (issue #27), and the
|
|
// Discord identity is the only thing that decides which Reader a login is: one
|
|
// code path serves the first login and every later one, so a returning Reader
|
|
// can never end up with a second library.
|
|
//
|
|
// epochZeroHash is only used for a brand-new row. An existing row keeps its
|
|
// stored hash untouched, or a login would silently undo a rotation and revive
|
|
// the credential the Reader rotated away from.
|
|
func (s *Store) EnsureReader(discordID string, epochZeroHash [32]byte) (int64, error) {
|
|
var id int64
|
|
// DO UPDATE rather than DO NOTHING because only an updated row is
|
|
// returned by RETURNING; assigning the column to itself is the no-op that
|
|
// makes the existing id come back.
|
|
err := s.db.QueryRow(`
|
|
INSERT INTO readers (discord_id, token_sha256) VALUES ($1, $2)
|
|
ON CONFLICT (discord_id) DO UPDATE SET discord_id = readers.discord_id
|
|
RETURNING id`, discordID, epochZeroHash[:]).Scan(&id)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("ensure reader: %w", err)
|
|
}
|
|
return id, nil
|
|
}
|
|
|
|
// SightingDisagreementLimit is the number of contradictions that stop that
|
|
// Reader's Sightings from deferring a Poll (issue #103). The counters exist
|
|
// before the mechanism that moves them, so the admin page (issue #102) can
|
|
// clear a false mark without waiting for the Sighting feature.
|
|
const SightingDisagreementLimit = 3
|
|
|
|
// Blocked reports whether this Reader's Sighting marks have reached the
|
|
// disagreement limit, which stops their Sightings from deferring a Poll.
|
|
func (r ReaderSummary) Blocked() bool {
|
|
return r.Disagreements >= SightingDisagreementLimit
|
|
}
|
|
|
|
// ReaderSummary is one Reader as the owner's administration panel sees them:
|
|
// who they are, how many live sessions they hold, and their Sighting marks.
|
|
// No credential material, hashed or otherwise, is exposed.
|
|
type ReaderSummary struct {
|
|
ID int64
|
|
DiscordID string
|
|
// Sessions counts unexpired session rows — what the owner revokes.
|
|
Sessions int
|
|
// Agreements and Disagreements are the Sighting counters (issue #102);
|
|
// zero means a trusted Reader.
|
|
Agreements int
|
|
Disagreements int
|
|
}
|
|
|
|
// Readers lists every Reader with their live session count and Sighting
|
|
// marks, oldest first, so the owner row (always the oldest) heads the list.
|
|
func (s *Store) Readers() ([]ReaderSummary, error) {
|
|
rows, err := s.db.Query(`
|
|
SELECT r.id, r.discord_id,
|
|
r.sighting_agreements, r.sighting_disagreements,
|
|
count(sess.id) FILTER (WHERE sess.expires_at > now()) AS sessions
|
|
FROM readers r
|
|
LEFT JOIN sessions sess ON sess.reader_id = r.id
|
|
GROUP BY r.id, r.discord_id, r.sighting_agreements, r.sighting_disagreements
|
|
ORDER BY r.id`)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("query readers: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
out := []ReaderSummary{}
|
|
for rows.Next() {
|
|
var r ReaderSummary
|
|
if err := rows.Scan(&r.ID, &r.DiscordID, &r.Agreements, &r.Disagreements, &r.Sessions); err != nil {
|
|
return nil, fmt.Errorf("scan reader: %w", err)
|
|
}
|
|
out = append(out, r)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// ClearReaderMarks zeroes a Reader's Sighting counters. It is the owner's
|
|
// remedy for a mark produced by a broken Site adapter rather than a dishonest
|
|
// Reader: it restores a privilege, it is not destruction.
|
|
func (s *Store) ClearReaderMarks(readerID int64) error {
|
|
if _, err := s.db.Exec(`
|
|
UPDATE readers SET sighting_agreements = 0, sighting_disagreements = 0
|
|
WHERE id = $1`, readerID); err != nil {
|
|
return fmt.Errorf("clear reader marks for reader %d: %w", readerID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// readersMigration is the version that creates the readers table. The owner
|
|
// seed runs between two migrate passes, so that the run-once migration which
|
|
// attaches existing bookmarks (0004) finds the owner row.
|
|
const readersMigration = 3
|
|
|
|
// allMigrations is the migrate() cap that applies every pending version.
|
|
const allMigrations = 0
|
|
|
|
// Open connects to Postgres at url — a libpq connection URL such as
|
|
// "postgres://user:pass@host:5432/bookmarks?sslmode=disable" — brings its
|
|
// schema up to date, seeds the owner Reader, and prepares cover storage.
|
|
func Open(url string, owner Owner, coverDir, coverBaseURL string) (*Store, error) {
|
|
if strings.TrimSpace(coverDir) == "" {
|
|
return nil, errors.New("cover directory is required")
|
|
}
|
|
// Every wire Cover is this string with a path glued on, rendered by a
|
|
// userscript on a Site's own origin: anything but an absolute origin
|
|
// produces addresses no client can load, silently (ADR-0007).
|
|
base := strings.TrimRight(coverBaseURL, "/")
|
|
if host, ok := strings.CutPrefix(base, "https://"); !ok || host == "" {
|
|
if host, ok := strings.CutPrefix(base, "http://"); !ok || host == "" {
|
|
return nil, fmt.Errorf("cover base URL %q is not an absolute http(s) origin", coverBaseURL)
|
|
}
|
|
}
|
|
if err := os.MkdirAll(coverDir, 0o755); err != nil {
|
|
return nil, fmt.Errorf("create cover directory: %w", err)
|
|
}
|
|
info, err := os.Stat(coverDir)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("stat cover directory: %w", err)
|
|
}
|
|
if !info.IsDir() {
|
|
return nil, fmt.Errorf("cover directory %q is not a directory", coverDir)
|
|
}
|
|
|
|
db, err := sql.Open("pgx", url)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("open postgres: %w", err)
|
|
}
|
|
// Schema runs in two passes with the seed between: 0003 creates the
|
|
// readers table, the owner row must exist before 0004 attaches the
|
|
// existing bookmarks to it. Anything past 0004 is applied by the second
|
|
// pass.
|
|
if err := migrate(db, readersMigration); err != nil {
|
|
db.Close()
|
|
return nil, fmt.Errorf("migrate schema: %w", err)
|
|
}
|
|
// The owner row must exist before 0004 attaches the existing bookmarks to
|
|
// it. The hash refresh is a separate statement after all migrations: the
|
|
// token_epoch column 0006 adds does not exist yet at this point, and the
|
|
// refresh only ever concerns rows that have never been rotated.
|
|
if err := seedOwner(db, owner); err != nil {
|
|
db.Close()
|
|
return nil, fmt.Errorf("seed owner: %w", err)
|
|
}
|
|
if err := migrate(db, allMigrations); err != nil {
|
|
db.Close()
|
|
return nil, fmt.Errorf("migrate: %w", err)
|
|
}
|
|
if err := refreshOwnerToken(db, owner); err != nil {
|
|
db.Close()
|
|
return nil, fmt.Errorf("refresh owner token: %w", err)
|
|
}
|
|
var ownerID int64
|
|
if err := db.QueryRow(
|
|
`SELECT id FROM readers WHERE discord_id = $1`, owner.DiscordID).Scan(&ownerID); err != nil {
|
|
db.Close()
|
|
return nil, fmt.Errorf("resolve owner: %w", err)
|
|
}
|
|
return &Store{
|
|
db: db, ownerID: ownerID, coverDir: coverDir, coverBaseURL: base,
|
|
}, nil
|
|
}
|
|
|
|
// seedOwner makes sure the configured owner exists as exactly one readers row.
|
|
// The hash is only ever written here for a brand-new row; existing rows keep
|
|
// what they have until refreshOwnerToken decides otherwise, so the seed can
|
|
// never clobber a rotation.
|
|
func seedOwner(db *sql.DB, o Owner) error {
|
|
if _, err := db.Exec(`
|
|
INSERT INTO readers (discord_id, token_sha256) VALUES ($1, $2)
|
|
ON CONFLICT (discord_id) DO NOTHING`,
|
|
o.DiscordID, o.TokenHash[:]); err != nil {
|
|
return fmt.Errorf("seed owner: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// refreshOwnerToken brings a never-rotated owner row's hash current with the
|
|
// configured credential. That is the cutover path: a database seeded under
|
|
// the retired global token still carries its hash at epoch 0, and the
|
|
// epoch-0 derivation is the caller's TokenHash. A rotated row (epoch > 0) is
|
|
// left alone — a restart must not resurrect the old credential by
|
|
// overwriting the hash a rotation wrote.
|
|
func refreshOwnerToken(db *sql.DB, o Owner) error {
|
|
if _, err := db.Exec(`
|
|
UPDATE readers SET token_sha256 = $2
|
|
WHERE discord_id = $1 AND token_epoch = 0`,
|
|
o.DiscordID, o.TokenHash[:]); err != nil {
|
|
return fmt.Errorf("refresh owner token: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// migrate applies every embedded migration this database has not recorded, in
|
|
// filename order, each in its own transaction. upto caps the highest version
|
|
// applied; 0 means all. Files are named "<version>_<name>.sql" and are
|
|
// append-only: editing an applied file changes nothing, because
|
|
// schema_migrations is how a database remembers what it ran. Runs on every
|
|
// start and is a no-op once current.
|
|
func migrate(db *sql.DB, upto int64) error {
|
|
if _, err := db.Exec(`CREATE TABLE IF NOT EXISTS schema_migrations (
|
|
version bigint PRIMARY KEY,
|
|
applied_at timestamptz NOT NULL DEFAULT now())`); err != nil {
|
|
return fmt.Errorf("create version table: %w", err)
|
|
}
|
|
|
|
names, err := fs.Glob(migrations, "migrations/*.sql")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
slices.Sort(names)
|
|
|
|
for _, name := range names {
|
|
version, err := strconv.ParseInt(strings.SplitN(path.Base(name), "_", 2)[0], 10, 64)
|
|
if err != nil {
|
|
return fmt.Errorf("migration %q: filename must start with a version number", name)
|
|
}
|
|
if upto > 0 && version > upto {
|
|
continue
|
|
}
|
|
body, err := migrations.ReadFile(name)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := applyMigration(db, version, string(body)); err != nil {
|
|
return fmt.Errorf("migration %q: %w", name, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// applyMigration runs one migration and records its version in the same
|
|
// transaction, so an interrupted start leaves neither half behind.
|
|
func applyMigration(db *sql.DB, version int64, body string) error {
|
|
tx, err := db.Begin()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer tx.Rollback()
|
|
|
|
var applied bool
|
|
if err := tx.QueryRow(
|
|
`SELECT EXISTS (SELECT 1 FROM schema_migrations WHERE version = $1)`,
|
|
version).Scan(&applied); err != nil {
|
|
return err
|
|
}
|
|
if applied {
|
|
return nil
|
|
}
|
|
// No parameters, so this goes over the simple protocol and a migration may
|
|
// hold more than one statement.
|
|
if _, err := tx.Exec(body); err != nil {
|
|
return err
|
|
}
|
|
if _, err := tx.Exec(`INSERT INTO schema_migrations (version) VALUES ($1)`, version); err != nil {
|
|
return err
|
|
}
|
|
return tx.Commit()
|
|
}
|
|
|
|
// scanBookmark reads one row in bookmarkColumns order. Every column is NOT
|
|
// NULL except latest_chapter_num, where NULL means "never captured" — a
|
|
// distinct state from chapter zero, and the reason for the pointer.
|
|
func (s *Store) scanBookmark(scan func(...any) error) (Bookmark, error) {
|
|
var (
|
|
b Bookmark
|
|
coverAddress string
|
|
latestChapterNum sql.NullFloat64
|
|
)
|
|
if err := scan(
|
|
&b.Site, &b.SeriesID, &b.Title, &b.SeriesURL, &coverAddress,
|
|
&b.LastChapter, &b.LastChapterNum, &b.LastChapterURL,
|
|
&b.Favorite, &b.LatestChapter, &latestChapterNum, &b.UpdatedAt, &b.Status, &b.Kind,
|
|
); err != nil {
|
|
return Bookmark{}, err
|
|
}
|
|
b.Cover = s.CoverWireURL(coverAddress)
|
|
if latestChapterNum.Valid {
|
|
b.LatestChapterNum = &latestChapterNum.Float64
|
|
}
|
|
// The wire identity is derived: there is no stored key column, the
|
|
// bookmark is keyed (reader_id, site, series_id) (issue #22).
|
|
b.Key = b.Site + ":" + b.SeriesID
|
|
// An unrecognised bucket (a hand-edited row) would leave the row in no list
|
|
// at all, so anything outside the three known buckets reads as the default
|
|
// rather than being passed through.
|
|
if b.Status != StatusReading && b.Status != StatusArchived && b.Status != StatusFinished {
|
|
b.Status = StatusReading
|
|
}
|
|
return b, nil
|
|
}
|
|
|
|
// scanSeries reads one row in seriesColumns order, plus the due query's
|
|
// forced flag and reader_count columns. latest_chapter_num and
|
|
// latest_raised_by are both nullable, same as latest_chapter_num on the
|
|
// bookmark read path.
|
|
func scanSeries(scan func(...any) error) (Series, error) {
|
|
var (
|
|
sr Series
|
|
latestChapterNum sql.NullFloat64
|
|
latestRaisedBy sql.NullInt64
|
|
)
|
|
if err := scan(
|
|
&sr.Site, &sr.SeriesID, &sr.Title, &sr.SeriesURL, &sr.Cover, &sr.CoverAddress,
|
|
&sr.Kind, &sr.LatestChapter, &latestChapterNum, &sr.LatestCheckedAt, &latestRaisedBy,
|
|
&sr.Forced, &sr.readerCount,
|
|
); err != nil {
|
|
return Series{}, err
|
|
}
|
|
if latestChapterNum.Valid {
|
|
sr.LatestChapterNum = &latestChapterNum.Float64
|
|
}
|
|
if latestRaisedBy.Valid {
|
|
sr.LatestRaisedBy = &latestRaisedBy.Int64
|
|
}
|
|
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() }
|
|
|
|
func coverSourceAddress(sourceURL string) string {
|
|
sum := sha256.Sum256([]byte(sourceURL))
|
|
return hex.EncodeToString(sum[:])
|
|
}
|
|
|
|
func coverRelativePath(address string) string {
|
|
return address[:2] + "/" + address[2:4] + "/" + address
|
|
}
|
|
|
|
func (s *Store) getCover(sourceURL string) ([]byte, string, bool, error) {
|
|
return s.getCoverByAddress(coverSourceAddress(sourceURL))
|
|
}
|
|
|
|
func (s *Store) getCoverByAddress(address string) ([]byte, string, bool, error) {
|
|
var relativePath, contentType string
|
|
err := s.db.QueryRow(
|
|
`SELECT path, content_type FROM covers WHERE address = $1`, address,
|
|
).Scan(&relativePath, &contentType)
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return nil, "", false, nil
|
|
}
|
|
if err != nil {
|
|
return nil, "", false, fmt.Errorf("get cover %q: %w", address, err)
|
|
}
|
|
expectedPath := coverRelativePath(address)
|
|
if relativePath != expectedPath {
|
|
return nil, "", false, fmt.Errorf("cover %q has unexpected path %q", address, relativePath)
|
|
}
|
|
body, err := os.ReadFile(filepath.Join(s.coverDir, filepath.FromSlash(relativePath)))
|
|
if errors.Is(err, fs.ErrNotExist) {
|
|
return nil, "", false, nil
|
|
}
|
|
if err != nil {
|
|
return nil, "", false, fmt.Errorf("read cover %q: %w", address, err)
|
|
}
|
|
return body, contentType, true, nil
|
|
}
|
|
|
|
func (s *Store) putCover(sourceURL string, body []byte, contentType string) (string, error) {
|
|
stored, ok := CoverContentType(contentType)
|
|
if !ok {
|
|
return "", fmt.Errorf("put cover %q: unsupported content type %q", sourceURL, contentType)
|
|
}
|
|
contentType = stored
|
|
address := CoverAddressForBytes(body)
|
|
relativePath := coverRelativePath(address)
|
|
coverPath := filepath.Join(s.coverDir, filepath.FromSlash(relativePath))
|
|
if err := os.MkdirAll(filepath.Dir(coverPath), 0o755); err != nil {
|
|
return "", fmt.Errorf("create cover shard: %w", err)
|
|
}
|
|
tmp, err := os.CreateTemp(filepath.Dir(coverPath), ".cover-*")
|
|
if err != nil {
|
|
return "", fmt.Errorf("create cover temp file: %w", err)
|
|
}
|
|
tmpName := tmp.Name()
|
|
defer os.Remove(tmpName)
|
|
if _, err := tmp.Write(body); err != nil {
|
|
tmp.Close()
|
|
return "", fmt.Errorf("write cover temp file: %w", err)
|
|
}
|
|
if err := tmp.Sync(); err != nil {
|
|
tmp.Close()
|
|
return "", fmt.Errorf("sync cover temp file: %w", err)
|
|
}
|
|
if err := tmp.Close(); err != nil {
|
|
return "", fmt.Errorf("close cover temp file: %w", err)
|
|
}
|
|
if err := os.Link(tmpName, coverPath); err != nil && !errors.Is(err, fs.ErrExist) {
|
|
return "", fmt.Errorf("install cover file: %w", err)
|
|
}
|
|
if _, err := s.db.Exec(`
|
|
INSERT INTO covers (address, path, content_type)
|
|
VALUES ($1, $2, $3)
|
|
ON CONFLICT (address) DO NOTHING`, address, relativePath, contentType); err != nil {
|
|
return "", fmt.Errorf("record cover %q: %w", address, err)
|
|
}
|
|
return address, nil
|
|
}
|
|
|
|
// ReclaimCover permanently removes a Cover nothing references: the sharded
|
|
// file first, the covers row last. A blank address is a no-op, and so is any
|
|
// address a Series row still points at — byte-identical artwork is one row by
|
|
// construction (ADR-0014), so reclaiming one Series' stranded bytes must not
|
|
// blank another's. The file goes first because the covers row is the handle:
|
|
// an interrupted run stays findable in SQL — covers rows unreferenced by any
|
|
// series cover_address — and re-running finishes the job, whereas deleting
|
|
// the row first would leave a file nothing names. A concurrent Forced Poll
|
|
// repointing a live Series at this address between the guard and the unlink
|
|
// is the repairable case: the missing file reads as ok=false and the next
|
|
// pass re-installs it. Failures are returned, never logged here — the caller
|
|
// logs and carries on — and a failed unlink leaves the row in place for a
|
|
// retry. A whole-table sweep, if ever wanted, is one SQL query over covers,
|
|
// not a tree walk and not this function.
|
|
func (s *Store) ReclaimCover(address string) error {
|
|
if address == "" {
|
|
return nil
|
|
}
|
|
var referenced int
|
|
err := s.db.QueryRow(`SELECT 1 FROM series WHERE cover_address = $1 LIMIT 1`, address).Scan(&referenced)
|
|
if err == nil {
|
|
return nil
|
|
}
|
|
if !errors.Is(err, sql.ErrNoRows) {
|
|
return fmt.Errorf("guard reclaim of cover %q: %w", address, err)
|
|
}
|
|
coverPath := filepath.Join(s.coverDir, filepath.FromSlash(coverRelativePath(address)))
|
|
if err := os.Remove(coverPath); err != nil && !errors.Is(err, fs.ErrNotExist) {
|
|
return fmt.Errorf("remove cover file %q: %w", address, err)
|
|
}
|
|
if _, err := s.db.Exec(`DELETE FROM covers WHERE address = $1`, address); err != nil {
|
|
return fmt.Errorf("delete cover row %q: %w", address, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// GetCover returns the immutable object a source URL's own hash names. Rows
|
|
// written before byte addressing (ADR-0014) are the only ones that ever reach
|
|
// it; it hashes the URL, so a byte-addressed Cover is invisible to it. Missing
|
|
// files are reported with ok=false so callers can retry acquisition later.
|
|
func (s *Store) GetCover(sourceURL string) ([]byte, string, bool, error) {
|
|
return s.getCover(sourceURL)
|
|
}
|
|
|
|
// PutCover persists bytes under their own content address (ADR-0014). A later
|
|
// write of the same bytes cannot replace the immutable object.
|
|
func (s *Store) PutCover(sourceURL string, body []byte, contentType string) error {
|
|
_, err := s.putCover(sourceURL, body, contentType)
|
|
return err
|
|
}
|
|
|
|
// CoverAddressForBytes is the content address body is stored under: the hex
|
|
// SHA-256 of the bytes, so identical artwork is one address and a re-art a
|
|
// new one. Legacy rows were addressed from their source URL instead and are
|
|
// never rehashed — both derivations coexist (ADR-0014).
|
|
func CoverAddressForBytes(body []byte) string {
|
|
sum := sha256.Sum256(body)
|
|
return hex.EncodeToString(sum[:])
|
|
}
|
|
|
|
// coverAddressRe is the shape of a stored address: 64 lowercase hex digits —
|
|
// the hex SHA-256 of the cover bytes, or of the source URL for legacy rows
|
|
// (ADR-0014). Request paths reach CoverByAddress, so the shape is checked
|
|
// before the value is ever turned into a filesystem path; byte-derived
|
|
// addresses keep the same shape, so the guard is unchanged.
|
|
var coverAddressRe = regexp.MustCompile(`^[0-9a-f]{64}$`)
|
|
|
|
// CoverByAddress returns the immutable object at one content address. An
|
|
// address that is not a stored one - malformed, unknown, or recorded but with
|
|
// its file gone - is reported with ok=false rather than as an error.
|
|
func (s *Store) CoverByAddress(address string) ([]byte, string, bool, error) {
|
|
if !coverAddressRe.MatchString(address) {
|
|
return nil, "", false, nil
|
|
}
|
|
return s.getCoverByAddress(address)
|
|
}
|
|
|
|
// CoverWireURL is the absolute URL a client renders for a stored Cover, and ""
|
|
// for a Series that has none yet. A blank is a real state, not a placeholder
|
|
// address: it is what tells both clients to draw their own fallback instead of
|
|
// requesting bytes that do not exist (ADR-0007).
|
|
func (s *Store) CoverWireURL(address string) string {
|
|
if address == "" {
|
|
return ""
|
|
}
|
|
return s.coverBaseURL + "/covers/" + address
|
|
}
|
|
|
|
// SetSeriesCover stores the bytes and points the Series at their address, but
|
|
// only while the Series has no Cover: acquisition at creation and the poll
|
|
// both call this, and whichever arrives second must not overwrite the first.
|
|
// The bytes themselves are content-addressed and immutable, so storing them
|
|
// twice is free. See ReplaceSeriesCover for the write that may move a Cover
|
|
// once one exists (ADR-0014).
|
|
func (s *Store) SetSeriesCover(site, seriesID, sourceURL string, body []byte, contentType string) error {
|
|
address, err := s.putCover(sourceURL, body, contentType)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if _, err := s.db.Exec(`
|
|
UPDATE series SET cover = $3, cover_address = $4
|
|
WHERE site = $1 AND series_id = $2 AND cover_address = ''`,
|
|
site, seriesID, sourceURL, address); err != nil {
|
|
return fmt.Errorf("set cover for %q: %w", site+":"+seriesID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// ReplaceSeriesCover stores the bytes and points the Series at their address
|
|
// whether or not one already exists, writing the current source URL alongside
|
|
// — the Forced Poll's installer and the only write that may move a Cover once
|
|
// one exists (ADR-0014). previous is the address the row held before the write
|
|
// ("" if it had none) and current the address of the bytes just stored; both
|
|
// are read and written in one transaction, so a concurrent replacement reports
|
|
// the exact displacement. previous == current means the Site served identical
|
|
// artwork, an honest no-op; otherwise previous is stranded — the row no
|
|
// longer points at it, and reclaiming its bytes is the caller's separate act
|
|
// (the poller's replace path calls ReclaimCover on it). This write itself
|
|
// removes nothing.
|
|
func (s *Store) ReplaceSeriesCover(site, seriesID, sourceURL string, body []byte, contentType string) (previous, current string, err error) {
|
|
current, err = s.putCover(sourceURL, body, contentType)
|
|
if err != nil {
|
|
return "", "", err
|
|
}
|
|
tx, err := s.db.Begin()
|
|
if err != nil {
|
|
return "", "", fmt.Errorf("begin replace cover for %q: %w", site+":"+seriesID, err)
|
|
}
|
|
defer tx.Rollback()
|
|
err = tx.QueryRow(`
|
|
SELECT cover_address FROM series
|
|
WHERE site = $1 AND series_id = $2 FOR UPDATE`,
|
|
site, seriesID).Scan(&previous)
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
previous = ""
|
|
} else if err != nil {
|
|
return "", "", fmt.Errorf("read cover for %q: %w", site+":"+seriesID, err)
|
|
}
|
|
if _, err := tx.Exec(`
|
|
UPDATE series SET cover = $3, cover_address = $4
|
|
WHERE site = $1 AND series_id = $2`,
|
|
site, seriesID, sourceURL, current); err != nil {
|
|
return "", "", fmt.Errorf("replace cover for %q: %w", site+":"+seriesID, err)
|
|
}
|
|
if err := tx.Commit(); err != nil {
|
|
return "", "", fmt.Errorf("commit cover replace for %q: %w", site+":"+seriesID, err)
|
|
}
|
|
return previous, current, nil
|
|
}
|
|
|
|
// List returns every bookmark of one reader, newest activity first.
|
|
// Series-owned fields are joined in, so each Bookmark reads back whole and
|
|
// flat (ADR-0004).
|
|
func (s *Store) List(readerID int64) ([]Bookmark, error) {
|
|
rows, err := s.db.Query(`SELECT `+bookmarkColumns+`
|
|
FROM bookmarks b
|
|
JOIN series s ON s.site = b.site AND s.series_id = b.series_id
|
|
WHERE b.reader_id = $1
|
|
ORDER BY b.updated_at DESC`, readerID)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("query bookmarks: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
out := []Bookmark{}
|
|
for rows.Next() {
|
|
b, err := s.scanBookmark(rows.Scan)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("scan bookmark: %w", err)
|
|
}
|
|
out = append(out, b)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// Get returns one bookmark of one reader by key. A missing key is not an
|
|
// error: ok is false and err is nil. UI mutations read-modify-write through
|
|
// this so they preserve the fields they do not touch.
|
|
func (s *Store) Get(readerID int64, key string) (Bookmark, bool, error) {
|
|
site, seriesID, ok := strings.Cut(key, ":")
|
|
if !ok {
|
|
return Bookmark{}, false, nil
|
|
}
|
|
b, err := s.scanBookmark(s.db.QueryRow(
|
|
`SELECT `+bookmarkColumns+` FROM bookmarks b
|
|
JOIN series s ON s.site = b.site AND s.series_id = b.series_id
|
|
WHERE b.reader_id = $1 AND b.site = $2 AND b.series_id = $3`,
|
|
readerID, site, seriesID).Scan)
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return Bookmark{}, false, nil
|
|
}
|
|
if err != nil {
|
|
return Bookmark{}, false, fmt.Errorf("get %q: %w", key, err)
|
|
}
|
|
return b, true, nil
|
|
}
|
|
|
|
// Upsert inserts or replaces one reader's bookmark by key (last-write-wins)
|
|
// and returns the row as actually stored — one flat object with the
|
|
// series-owned fields joined in, exactly as GET reports it (ADR-0004). A
|
|
// bookmark is keyed (reader_id, site, series_id), so the same key upserts two
|
|
// independent rows for two readers.
|
|
//
|
|
// The flat body is decomposed across two tables in one transaction. The series
|
|
// row is written first (the bookmarks FK requires it to exist), then the
|
|
// bookmark row. On the series side, title/series_url/cover are applied only
|
|
// when the row is brand new: once a series exists, client-supplied values are
|
|
// ignored, because the row is shared and the values are scraped page content —
|
|
// see ADR-0003. Kind and the latest-chapter fields are last-write-wins.
|
|
//
|
|
// b.UpdatedAt is only a candidate: it is applied when the row is new or when
|
|
// last_chapter_num changes, and otherwise the stored value is kept. Clients
|
|
// order their list by updated_at, so favoriting a series or recording a newly
|
|
// published chapter must not disturb that order — only real reading progress
|
|
// does. Callers must therefore use the returned bookmark, not the argument.
|
|
func (s *Store) Upsert(readerID int64, b Bookmark) (Bookmark, error) {
|
|
tx, err := s.db.Begin()
|
|
if err != nil {
|
|
return Bookmark{}, fmt.Errorf("begin %q: %w", b.Key, err)
|
|
}
|
|
defer tx.Rollback()
|
|
|
|
var latestNum any
|
|
if b.LatestChapterNum != nil {
|
|
latestNum = *b.LatestChapterNum
|
|
}
|
|
|
|
// The kind column resolves on the VALUES side, not in the conflict clause:
|
|
// excluded.* is the row *after* these expressions are evaluated, so a
|
|
// default applied there would look identical to a real 'manga' and would
|
|
// overwrite a novel series on every PUT from a client that knows nothing
|
|
// about the column. Resolved once here, an empty incoming kind means "keep
|
|
// what is stored", and only a brand-new row falls through to the literal
|
|
// default. The subquery runs inside this transaction, so it sees the row
|
|
// this statement is about to conflict with. Same pattern as the status
|
|
// COALESCE on the bookmark insert below.
|
|
//
|
|
// The ::text casts are load-bearing: inside COALESCE/NULLIF there is no
|
|
// target column to infer the parameter type from, and Postgres rejects the
|
|
// statement rather than guessing.
|
|
//
|
|
// The cover columns are absent on purpose: the Cover is acquired
|
|
// server-side (ADR-0007), so a client-supplied one is not written even
|
|
// when the row is brand new.
|
|
//
|
|
// xmax is zero only on a row this statement inserted, which is how a
|
|
// Series nobody had bookmarked before is told apart from one that already
|
|
// existed — DO UPDATE returns a row either way.
|
|
// latest_corrected_at is the one clause conditional on the value moving
|
|
// (#149): after a Correction a Reader's cached row holds the corrected
|
|
// number and resends it on the next Progress PUT, so unconditional
|
|
// zeroing would erase the fact while the value is still the owner's. The
|
|
// stamp survives a same-number PUT and dies the moment the number moves.
|
|
var created bool
|
|
if err := tx.QueryRow(`
|
|
INSERT INTO series (site, series_id, title, series_url, kind,
|
|
latest_chapter, latest_chapter_num)
|
|
VALUES ($1, $2, $3, $4,
|
|
COALESCE(NULLIF($5::text, ''), (SELECT kind FROM series WHERE site = $1 AND series_id = $2), 'manga'),
|
|
$6, $7)
|
|
ON CONFLICT (site, series_id) DO UPDATE SET
|
|
kind=excluded.kind,
|
|
latest_chapter=excluded.latest_chapter,
|
|
latest_chapter_num=excluded.latest_chapter_num,
|
|
latest_corrected_at = CASE
|
|
WHEN series.latest_chapter_num IS DISTINCT FROM excluded.latest_chapter_num
|
|
THEN 0 ELSE series.latest_corrected_at END
|
|
RETURNING xmax = 0`,
|
|
b.Site, b.SeriesID, b.Title, b.SeriesURL, b.Kind,
|
|
b.LatestChapter, latestNum).Scan(&created); err != nil {
|
|
return Bookmark{}, fmt.Errorf("upsert series for %q: %w", b.Key, err)
|
|
}
|
|
|
|
// IS DISTINCT FROM is Postgres's null-safe comparison, and it is what
|
|
// implements the ordering rule. Within DO UPDATE, a bare column is the
|
|
// stored row and excluded.* is the incoming one; a brand-new key never
|
|
// reaches this clause, so it keeps the fresh timestamp from VALUES.
|
|
if _, err := tx.Exec(`
|
|
INSERT INTO bookmarks (reader_id, site, series_id, last_chapter, last_chapter_num,
|
|
last_chapter_url, favorite, status, updated_at)
|
|
VALUES ($1, $2, $3, $4, $5, $6, $7,
|
|
COALESCE(NULLIF($8::text, ''), (SELECT status FROM bookmarks WHERE reader_id = $1 AND site = $2 AND series_id = $3), 'reading'),
|
|
$9)
|
|
ON CONFLICT (reader_id, site, series_id) DO UPDATE SET
|
|
last_chapter=excluded.last_chapter, last_chapter_num=excluded.last_chapter_num,
|
|
last_chapter_url=excluded.last_chapter_url,
|
|
favorite=excluded.favorite,
|
|
status=excluded.status,
|
|
updated_at=CASE
|
|
WHEN bookmarks.last_chapter_num IS DISTINCT FROM excluded.last_chapter_num
|
|
THEN excluded.updated_at
|
|
ELSE bookmarks.updated_at
|
|
END`,
|
|
readerID, b.Site, b.SeriesID,
|
|
b.LastChapter, b.LastChapterNum, b.LastChapterURL,
|
|
b.Favorite, b.Status, b.UpdatedAt); err != nil {
|
|
return Bookmark{}, fmt.Errorf("upsert %q: %w", b.Key, err)
|
|
}
|
|
|
|
stored, err := s.scanBookmark(tx.QueryRow(
|
|
`SELECT `+bookmarkColumns+` FROM bookmarks b
|
|
JOIN series s ON s.site = b.site AND s.series_id = b.series_id
|
|
WHERE b.reader_id = $1 AND b.site = $2 AND b.series_id = $3`,
|
|
readerID, b.Site, b.SeriesID).Scan)
|
|
if err != nil {
|
|
return Bookmark{}, fmt.Errorf("read back %q: %w", b.Key, err)
|
|
}
|
|
if err := tx.Commit(); err != nil {
|
|
return Bookmark{}, fmt.Errorf("commit %q: %w", b.Key, err)
|
|
}
|
|
// After commit, never inside the transaction: the hook reaches a
|
|
// third-party Site, and the Reader's write must not wait on it.
|
|
if created && s.OnSeriesCreated != nil {
|
|
s.OnSeriesCreated(Series{
|
|
Site: b.Site, SeriesID: b.SeriesID, Title: stored.Title,
|
|
SeriesURL: stored.SeriesURL, Kind: stored.Kind,
|
|
})
|
|
}
|
|
return stored, nil
|
|
}
|
|
|
|
// Delete removes one reader's bookmark by key. Deleting a missing key is not
|
|
// an error.
|
|
func (s *Store) Delete(readerID int64, key string) error {
|
|
site, seriesID, ok := strings.Cut(key, ":")
|
|
if !ok {
|
|
return nil
|
|
}
|
|
if _, err := s.db.Exec(
|
|
`DELETE FROM bookmarks WHERE reader_id = $1 AND site = $2 AND series_id = $3`,
|
|
readerID, site, seriesID); err != nil {
|
|
return fmt.Errorf("delete %q: %w", key, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// pgForeignKeyViolation is the SQLSTATE the driver surfaces when a Bookmark
|
|
// row refuses a Series delete (bookmarks_series_fk). pgconn exports no named
|
|
// constant for it, so the store names it here.
|
|
const pgForeignKeyViolation = "23503"
|
|
|
|
// ErrSeriesHasBookmarks is RemoveSeries' refusal: a Reader still holds the
|
|
// Series, so the owner's removal must not reach past that record. The
|
|
// delete is the check — no NOT EXISTS pre-check that can race the insert —
|
|
// and the driver's foreign-key violation is translated here so no driver
|
|
// type escapes the store (issue #155).
|
|
var ErrSeriesHasBookmarks = errors.New("series has bookmarks")
|
|
|
|
// RemoveSeries deletes one Series row by (site, series_id). It is refused
|
|
// while any Bookmark references the row; deleting an absent key is not an
|
|
// error, matching Delete. The caller owns the stranded Cover: read the row's
|
|
// cover_address before the delete and call ReclaimCover after it — the
|
|
// helper's guard cannot pass while the series row still points at the
|
|
// address, so the order is the sequence, not a preference.
|
|
func (s *Store) RemoveSeries(site, seriesID string) error {
|
|
if _, err := s.db.Exec(
|
|
`DELETE FROM series WHERE site = $1 AND series_id = $2`,
|
|
site, seriesID); err != nil {
|
|
var pgErr *pgconn.PgError
|
|
if errors.As(err, &pgErr) && pgErr.Code == pgForeignKeyViolation {
|
|
return ErrSeriesHasBookmarks
|
|
}
|
|
return fmt.Errorf("remove series %s:%s: %w", site, seriesID, err)
|
|
}
|
|
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()
|
|
}
|
|
|
|
// LaneGates reads a Site's pause and refusal stamps in one row read — the
|
|
// top-of-pass gate the poller uses (issue #141). A missing state row is the
|
|
// default: unpaused and not refusing.
|
|
func (s *Store) LaneGates(site string) (pausedUntil, refuseUntil int64, err error) {
|
|
err = s.db.QueryRow(
|
|
`SELECT paused_until, refuse_until FROM poll_lanes WHERE site = $1`, site).
|
|
Scan(&pausedUntil, &refuseUntil)
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return 0, 0, nil
|
|
}
|
|
if err != nil {
|
|
return 0, 0, fmt.Errorf("lane gates %s: %w", site, err)
|
|
}
|
|
return pausedUntil, refuseUntil, 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
|
|
// query, because each Poll Lane asks for its own list: the query carries one
|
|
// Site and one cut-off instead of parallel lists (issue #100). There is no
|
|
// limit — the Lane's own gap paces the fetches, and the batch size that used
|
|
// to cap this query is gone with the shared pace.
|
|
//
|
|
// A forced Series (force_poll_at newer than latest_checked_at, issue #146)
|
|
// overrides exactly three gates: the rest cutoff, the Sighting-deferral
|
|
// clause and the finished-only bucket. It never overrides an empty
|
|
// series_url or the Bookmarks join — nothing to fetch, and no consumer for
|
|
// the result — so those stay unconditional. Forced rows sort to the front of
|
|
// the queue; the reader-count-then-age ordering among the rest is ADR-0003.
|
|
//
|
|
// The reader_count ordering is the point of the split (ADR-0003): a series
|
|
// shared by several readers is fetched once per due cycle, and the popular
|
|
// ones stay freshest while the long tail absorbs any shortfall. Within one
|
|
// reader count, oldest-first keeps the poll fair when the backlog outgrows
|
|
// throughput: the most neglected series is always next, so a large collection
|
|
// refreshes uniformly slower rather than leaving a tail that never refreshes
|
|
// at all. The userscript sorts its own queue the same way (L453).
|
|
//
|
|
// Series with no series_url are skipped — there is nothing to fetch, which is
|
|
// the same filter the userscript applies at L452. Series whose only bookmarks
|
|
// are finished are skipped too: nothing more is coming, so fetching them only
|
|
// burns requests. Archived bookmarks still count — knowing what a shelved
|
|
// series is up to is the whole reason for archiving instead of deleting.
|
|
// A series with no bookmarks at all never appears: the join excludes it.
|
|
//
|
|
// ceilingMs is the Sighting deferral ceiling (issue #103): a Series whose last
|
|
// real Poll is older than it appears however recently it was sighted. That is
|
|
// what bounds the whole mechanism — a wrong Latest Chapter dies within the
|
|
// ceiling deterministically rather than in expectation. Deferral itself is
|
|
// decided here, from two facts the query already computes, so a Lane gains no
|
|
// query per round: a Sighting younger than cutoffMs holds the Series back, but
|
|
// only while COUNT(*) is 1. A Series a second Reader bookmarks is Polled on
|
|
// schedule, so a wrong value the whole guild can see is corrected by a check
|
|
// that was never postponed; on a solitary Series the only person a wrong value
|
|
// reaches is the Reader who reported it. Whether the reporting Reader is
|
|
// allowed to defer at all was settled when the Sighting was recorded — see
|
|
// RecordSighting.
|
|
func (s *Store) DueForLatestCheck(site string, cutoffMs, ceilingMs int64) ([]Series, error) {
|
|
rows, err := s.db.Query(`SELECT `+seriesColumns+`,
|
|
(s.force_poll_at > s.latest_checked_at) AS forced,
|
|
COUNT(*) AS reader_count
|
|
FROM series s
|
|
JOIN bookmarks b ON b.site = s.site AND b.series_id = s.series_id
|
|
WHERE s.site = $1
|
|
AND s.series_url <> ''
|
|
AND (s.latest_checked_at <= $2::bigint
|
|
OR s.force_poll_at > s.latest_checked_at)
|
|
GROUP BY s.site, s.series_id, s.title, s.series_url, s.cover,
|
|
s.kind, s.latest_chapter, s.latest_chapter_num, s.latest_checked_at,
|
|
s.force_poll_at
|
|
HAVING (COUNT(*) FILTER (WHERE b.status <> 'finished') > 0
|
|
OR s.force_poll_at > s.latest_checked_at)
|
|
AND (COUNT(*) > 1
|
|
OR s.latest_sighted_at <= $2::bigint
|
|
OR s.latest_checked_at <= $3::bigint
|
|
OR s.force_poll_at > s.latest_checked_at)
|
|
ORDER BY (s.force_poll_at > s.latest_checked_at) DESC,
|
|
reader_count DESC, s.latest_checked_at ASC`, site, cutoffMs, ceilingMs)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("query due series: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
out := []Series{}
|
|
for rows.Next() {
|
|
sr, err := scanSeries(rows.Scan)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("scan due series: %w", err)
|
|
}
|
|
out = append(out, sr)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// EligibleSeriesCount returns how many of a Site's Series still have at least
|
|
// one bookmark outside the finished bucket. It is the denominator of the
|
|
// Lane's pace (issue #100): the effective gap is the smaller of the registry
|
|
// gap and one hour divided by this count, so Series that will never be Polled
|
|
// do not make the Lane faster than it needs to be, and counting every eligible
|
|
// Series rather than only those currently due keeps the pace steady — the
|
|
// single worst moment to be fastest is startup, when everything is due at
|
|
// once.
|
|
func (s *Store) EligibleSeriesCount(site string) (int, error) {
|
|
var n int
|
|
err := s.db.QueryRow(`SELECT COUNT(*) FROM (
|
|
SELECT 1
|
|
FROM series s
|
|
JOIN bookmarks b ON b.site = s.site AND b.series_id = s.series_id
|
|
WHERE s.site = $1
|
|
GROUP BY s.site, s.series_id
|
|
HAVING COUNT(*) FILTER (WHERE b.status <> 'finished') > 0
|
|
) e`, site).Scan(&n)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("count eligible series %s: %w", site, err)
|
|
}
|
|
return n, nil
|
|
}
|
|
|
|
// MarkLatestChecked records that the server looked at a series at ts, whatever
|
|
// the look turned up. Marking a missing series is not an error: the row may
|
|
// have been orphaned while a fetch was in flight.
|
|
//
|
|
// This is the one write that does not go through Upsert, and the column is kept
|
|
// out of the client-visible read path on purpose. PUT /bookmarks/{key} decodes
|
|
// a whole Bookmark from the client and Upsert writes every series column it
|
|
// knows about, so a userscript PUT — which has no idea this field exists —
|
|
// would write a zero and reset the rest, making the poller re-fetch that
|
|
// series every tick for as long as the user kept reading it.
|
|
func (s *Store) MarkLatestChecked(site, seriesID string, ts int64) error {
|
|
if _, err := s.db.Exec(
|
|
`UPDATE series SET latest_checked_at = $1 WHERE site = $2 AND series_id = $3`,
|
|
ts, site, seriesID); err != nil {
|
|
return fmt.Errorf("mark checked %s:%s: %w", site, seriesID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// ForceSeriesPoll stamps a Series with the owner's "check now" request
|
|
// (issue #146): a fact about the Series the Lane's next pass reads through
|
|
// DueForLatestCheck, never a command to the poller — so the request survives
|
|
// a restart. Writing again overwrites the request time; the write is
|
|
// idempotent. Touching a missing series is not an error: the row may have
|
|
// been orphaned, and the caller's read decides what exists. The stamp never
|
|
// expires by itself — an unanswered request keeps ageing — and pending is
|
|
// derived as force_poll_at > latest_checked_at, which is why the poller's
|
|
// check stamp is written before the fetch: the first attempt ends the
|
|
// pending state whatever it returns.
|
|
func (s *Store) ForceSeriesPoll(site, seriesID string, at int64) error {
|
|
if _, err := s.db.Exec(
|
|
`UPDATE series SET force_poll_at = $1 WHERE site = $2 AND series_id = $3`,
|
|
at, site, seriesID); err != nil {
|
|
return fmt.Errorf("force poll %s:%s: %w", site, seriesID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// LatestCheckedAt reads the column MarkLatestChecked writes. It exists for
|
|
// tests outside this package (the poller's own tests assert on rest
|
|
// bookkeeping) — see MarkLatestChecked for why the field stays off the
|
|
// client-visible row.
|
|
func (s *Store) LatestCheckedAt(site, seriesID string) (int64, error) {
|
|
var ts int64
|
|
if err := s.db.QueryRow(
|
|
`SELECT latest_checked_at FROM series WHERE site = $1 AND series_id = $2`,
|
|
site, seriesID).Scan(&ts); err != nil {
|
|
return 0, fmt.Errorf("latest checked at %s:%s: %w", site, seriesID, err)
|
|
}
|
|
return ts, nil
|
|
}
|
|
|
|
// SetLatestChapter records the newest chapter the poll found on a series page.
|
|
// The poller walks Series rather than Bookmarks, so this is a series-level
|
|
// write: the row is shared, and updating it once refreshes every bookmark that
|
|
// joins to it. Touching a missing series is not an error. The correction stamp
|
|
// is zeroed unconditionally: checkOne only calls this when the number differs,
|
|
// so a second copy of the condition would drift (#149).
|
|
func (s *Store) SetLatestChapter(site, seriesID, label string, num float64) error {
|
|
if _, err := s.db.Exec(
|
|
`UPDATE series SET latest_chapter = $3, latest_chapter_num = $4,
|
|
latest_corrected_at = 0
|
|
WHERE site = $1 AND series_id = $2`,
|
|
site, seriesID, label, num); err != nil {
|
|
return fmt.Errorf("set latest chapter %s:%s: %w", site, seriesID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SetSeriesURL stores the owner's repair for a Series' source address
|
|
// (issue #151): the one write that lifts the write-once rule documented on
|
|
// Series.SeriesURL. It is a store, not a verification — the caller has
|
|
// already passed the poller's fetch gate. The handler 404s on an unknown row
|
|
// before calling; the write itself is a plain single-column UPDATE like
|
|
// MarkLatestChecked.
|
|
func (s *Store) SetSeriesURL(site, seriesID, seriesURL string) error {
|
|
if _, err := s.db.Exec(
|
|
`UPDATE series SET series_url = $3 WHERE site = $1 AND series_id = $2`,
|
|
site, seriesID, seriesURL); err != nil {
|
|
return fmt.Errorf("set series url %s:%s: %w", site, seriesID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// CorrectLatestChapter makes the Latest Chapter the owner's: one UPDATE
|
|
// carrying the number, the derived label and the correction stamp. The label
|
|
// shape is the poller's and the userscript's ("Chapter " + the number as
|
|
// printed), so chapterLeadIn strips it and the UI renders "Ch N" with no
|
|
// special case. latest_checked_at is not touched: a Correction is not a check.
|
|
// A raising Reader is cleared without judgement: the number is the owner's
|
|
// now, and no Sighting counter moves (spec #135).
|
|
func (s *Store) CorrectLatestChapter(site, seriesID string, num float64, at int64) error {
|
|
if _, err := s.db.Exec(
|
|
`UPDATE series SET
|
|
latest_chapter = $3,
|
|
latest_chapter_num = $4,
|
|
latest_corrected_at = $5,
|
|
latest_raised_by = NULL
|
|
WHERE site = $1 AND series_id = $2`,
|
|
site, seriesID, "Chapter "+strconv.FormatFloat(num, 'f', -1, 64), num, at); err != nil {
|
|
return fmt.Errorf("correct latest chapter %s:%s: %w", site, seriesID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// (issue #103). It must be called *before* the Upsert that stores the reported
|
|
// value: the raise test compares against what is still on the row, and after
|
|
// the Upsert there is nothing left to compare with. A Series that does not
|
|
// exist yet — the first Bookmark of it — is not a Sighting at all: nothing has
|
|
// ever been Polled, so there is nothing to defer and nobody to attribute.
|
|
//
|
|
// Two independent effects, hence the two CASE arms. The deferral stamp is only
|
|
// written for a Reader below the disagreement limit, so a marked Reader's
|
|
// reports keep updating the Latest Chapter but stop postponing anything, and
|
|
// clearing their marks restores the privilege on their next Sighting. The
|
|
// attribution is written whenever the report raises the stored number,
|
|
// including for a marked Reader — their Sightings are still judged, which is
|
|
// how they earn the privilege back.
|
|
//
|
|
// num is the reported chapter number. A PUT that carries none — a favourite
|
|
// toggle, or progress written from a chapter page — is no Sighting at all:
|
|
// nobody read the Series page, so there is nothing to stand in for a Poll and
|
|
// nothing that could later be judged.
|
|
func (s *Store) RecordSighting(readerID int64, site, seriesID string, num *float64, ts int64) error {
|
|
if num == nil {
|
|
return nil
|
|
}
|
|
if _, err := s.db.Exec(`
|
|
UPDATE series SET
|
|
latest_sighted_at = CASE
|
|
WHEN (SELECT sighting_disagreements FROM readers WHERE id = $3) < $6
|
|
THEN $4::bigint ELSE latest_sighted_at END,
|
|
latest_raised_by = CASE
|
|
WHEN latest_chapter_num IS NULL OR $5::double precision > latest_chapter_num
|
|
THEN $3::bigint ELSE latest_raised_by END
|
|
WHERE site = $1 AND series_id = $2`,
|
|
site, seriesID, readerID, ts, *num, SightingDisagreementLimit); err != nil {
|
|
return fmt.Errorf("record sighting %s:%s: %w", site, seriesID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SightingAgreementsToClear is how many Polls must confirm a Reader's
|
|
// Sightings in a row before their disagreements are forgiven. An agreement is
|
|
// only recorded when a Poll later confirms a Sighting, so this is twenty Polls
|
|
// of Series that Reader bookmarks — hours to days, not twenty page views. That
|
|
// is the intended price: recovery is automatic but cannot be outwaited, and a
|
|
// disagreement resets the run to zero, so credit cannot be banked in advance.
|
|
const SightingAgreementsToClear = 20
|
|
|
|
// RecordSightingOutcome settles what a Poll decided about the Reader whose
|
|
// Sighting last raised this Series' Latest Chapter, and clears the attribution
|
|
// in the same transaction so one Sighting is judged exactly once. agreed is
|
|
// the Poll confirming the stored value; its opposite is the Poll finding a
|
|
// lower number, which means the raise was false.
|
|
//
|
|
// A Poll finding a *higher* number is neither — the Site published — and takes
|
|
// ClearSightingAttribution instead.
|
|
func (s *Store) RecordSightingOutcome(site, seriesID string, readerID int64, agreed bool) error {
|
|
tx, err := s.db.Begin()
|
|
if err != nil {
|
|
return fmt.Errorf("begin sighting outcome %s:%s: %w", site, seriesID, err)
|
|
}
|
|
defer tx.Rollback()
|
|
|
|
// The run length is what "consecutive" means: a disagreement zeroes the
|
|
// agreements, and completing a run zeroes both, so the next run starts
|
|
// from nothing rather than forgiving every later disagreement instantly.
|
|
q := `UPDATE readers SET sighting_disagreements = sighting_disagreements + 1,
|
|
sighting_agreements = 0
|
|
WHERE id = $1`
|
|
args := []any{readerID}
|
|
if agreed {
|
|
q = `UPDATE readers SET
|
|
sighting_agreements = CASE WHEN sighting_agreements + 1 >= $2 THEN 0
|
|
ELSE sighting_agreements + 1 END,
|
|
sighting_disagreements = CASE WHEN sighting_agreements + 1 >= $2 THEN 0
|
|
ELSE sighting_disagreements END
|
|
WHERE id = $1`
|
|
args = append(args, SightingAgreementsToClear)
|
|
}
|
|
if _, err := tx.Exec(q, args...); err != nil {
|
|
return fmt.Errorf("record sighting outcome for reader %d: %w", readerID, err)
|
|
}
|
|
if _, err := tx.Exec(clearAttributionSQL, site, seriesID, readerID); err != nil {
|
|
return fmt.Errorf("clear sighting attribution %s:%s: %w", site, seriesID, err)
|
|
}
|
|
if err := tx.Commit(); err != nil {
|
|
return fmt.Errorf("commit sighting outcome %s:%s: %w", site, seriesID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// ClearSightingAttribution answers a Sighting without judging it: the Poll
|
|
// found a higher number, so the value about to be stored is its own and this
|
|
// Reader is no longer answerable for the row. Without it the next Poll's
|
|
// agreement would be credited to a Reader who did not earn it.
|
|
func (s *Store) ClearSightingAttribution(site, seriesID string, readerID int64) error {
|
|
if _, err := s.db.Exec(clearAttributionSQL, site, seriesID, readerID); err != nil {
|
|
return fmt.Errorf("clear sighting attribution %s:%s: %w", site, seriesID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// clearAttributionSQL drops the attribution only while it still names the
|
|
// Reader being judged: a Sighting landing between the due query's snapshot and
|
|
// this write is a fresh, unjudged one and must not be erased by the previous
|
|
// one's verdict.
|
|
const clearAttributionSQL = `UPDATE series SET latest_raised_by = NULL
|
|
WHERE site = $1 AND series_id = $2 AND latest_raised_by = $3`
|