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 ":" 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, NotFound 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, NotFound int } // LanePause is one persisted Lane pause stamp. type LanePause struct { Site string PausedUntil int64 } // Key returns the canonical identity in bookmark-key form (":"), // 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, p.not_found, 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 "_.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.NotFound, &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, not_found) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13)`, 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.NotFound); 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, not_found 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), SUM(not_found) 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, &outcomes.NotFound, ); 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`