8e4fa6448e
Closes #136. Spec #136 end to end: `finished` becomes a fact about the Series, written only by the owner, and the reader-facing Lifecycle bucket is gone. ## What landed - **#157** — `series.finished_at bigint NOT NULL DEFAULT 0` plus the migration whose statement order is load-bearing (seed from the buckets, then flip them); both Lane queries lose the `HAVING COUNT(*) FILTER (WHERE b.status <> 'finished')` clause and gate on `finished_at = 0` instead, with the due-query/eligible-count force asymmetry kept deliberate and commented; `StatusFinished`, its API special-case 400, the web tab and the templates' Finished bucket deleted. - **#158** — owner Finish control on the Series detail page: confirm-gated finish, instant un-finish, admin accent (never ember, nothing is destroyed), `Store.SetSeriesFinished`, the two routes behind the owner gate, and the state displayed on the list row without offering the control there. - **#160** — reader side: derived `finished` bool on the flat Bookmark (`s.finished_at > 0`), rendered as a text-only label in both userscripts and on the web card; read-only inbound by omission from `Upsert`'s explicit `series` column list, same mechanism that already protects `cover`. - **#161** — glossary and the stale Reader-count divergence note catch up. - **#159** — `finished` joins the admin filter vocabulary (predicate `finished_at > 0`, label `Finished`, own aggregate count, figure last in the stats block as informational); the four clock-driven hygiene predicates (stale, never-checked, no-cover, no-chapter) exclude finished Series while unpollable, orphan and sighting-raised deliberately do not. ## Verification `go vet ./...` and `go test ./...` green on the merged branch (Docker-backed, throwaway `postgres:17-alpine` per package). Each ticket also passed a two-axis review (spec + standards) on its own branch before merge. Reviewed-on: #163 Co-authored-by: Sulthan Zaki <sultankiki05@gmail.com> Co-committed-by: Sulthan Zaki <sultankiki05@gmail.com>
1635 lines
67 KiB
Go
1635 lines
67 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 or archived.
|
|
// Archived series stay polled for new chapters.
|
|
// 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"`
|
|
// Finished is the owner's retirement of the Series, derived: the flag is a
|
|
// Series fact and a client cannot write it — see Upsert.
|
|
Finished bool `json:"finished"`
|
|
}
|
|
|
|
// 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.
|
|
// Finished is not a bucket: it is a fact about the Series (series.finished_at),
|
|
// never about a Reader's bookmark.
|
|
const (
|
|
StatusReading = "reading"
|
|
StatusArchived = "archived"
|
|
)
|
|
|
|
//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,
|
|
s.finished_at > 0`
|
|
|
|
// 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,
|
|
&b.Finished,
|
|
); 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 two known buckets reads as the default
|
|
// rather than being passed through.
|
|
if b.Status != StatusReading && b.Status != StatusArchived {
|
|
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 a finished Series. It never overrides an empty series_url or the
|
|
// Bookmarks join — nothing to fetch, and no consumer for the result — so
|
|
// those stay unconditional, and it never clears the finish: nothing here
|
|
// writes finished_at, and pending force clears itself when the pass stamps
|
|
// the check timestamp. 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.
|
|
//
|
|
// Series with no series_url are skipped — there is nothing to fetch, which is
|
|
// the same filter the userscript applies before refreshing. A finished Series is
|
|
// skipped unless forced: nothing more is coming, so fetching it only burns
|
|
// requests (issue #157). 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
|
|
// the schedule, so a wrong value the 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)
|
|
AND (s.finished_at = 0 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, s.finished_at
|
|
HAVING (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 are not finished —
|
|
// the flag, never a Reader vote. 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.
|
|
//
|
|
// The deliberate asymmetry with DueForLatestCheck's WHERE: a forced Series
|
|
// is due but never admitted here, because a forced pass must not speed up
|
|
// every other fetch on the Site — one impassioned press is not a reason to
|
|
// hammer the Site (issue #157).
|
|
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
|
|
AND s.finished_at = 0
|
|
GROUP BY s.site, s.series_id
|
|
) 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
|
|
}
|
|
|
|
// SetSeriesFinished stamps or clears the owner's finish. at is unix ms to
|
|
// finish, zero to un-finish. A finished Series drops out of the Lane's reads
|
|
// (issue #157), and nothing else writes this column: it is the only writer
|
|
// outside migration 0016, so a machine write can never retire a Series
|
|
// silently. Touching a missing series is not an error: the row may have been
|
|
// orphaned, and the caller's read decides what exists.
|
|
func (s *Store) SetSeriesFinished(site, seriesID string, at int64) error {
|
|
if _, err := s.db.Exec(
|
|
`UPDATE series SET finished_at = $1 WHERE site = $2 AND series_id = $3`,
|
|
at, site, seriesID); err != nil {
|
|
return fmt.Errorf("set series finished %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`
|