Move cover bytes to filesystem storage
This commit is contained in:
@@ -0,0 +1,9 @@
|
||||
-- Cover bytes move out of Postgres. Existing rows are intentionally dropped:
|
||||
-- the old kagane path already refetches missing Covers on demand.
|
||||
DROP TABLE covers;
|
||||
|
||||
CREATE TABLE covers (
|
||||
address text PRIMARY KEY,
|
||||
path text NOT NULL,
|
||||
content_type text NOT NULL
|
||||
);
|
||||
@@ -1,12 +1,16 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"database/sql"
|
||||
"embed"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io/fs"
|
||||
"os"
|
||||
"path"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
"slices"
|
||||
"strconv"
|
||||
@@ -225,7 +229,8 @@ type Store struct {
|
||||
// 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
|
||||
ownerID int64
|
||||
coverDir string
|
||||
}
|
||||
|
||||
// OwnerID returns the seeded owner Reader's id: the administrator, and the
|
||||
@@ -360,8 +365,22 @@ 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, and seeds the owner Reader.
|
||||
func Open(url string, owner Owner) (*Store, error) {
|
||||
// schema up to date, seeds the owner Reader, and prepares cover storage.
|
||||
func Open(url string, owner Owner, coverDir string) (*Store, error) {
|
||||
if strings.TrimSpace(coverDir) == "" {
|
||||
return nil, errors.New("cover directory is required")
|
||||
}
|
||||
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)
|
||||
@@ -396,7 +415,7 @@ func Open(url string, owner Owner) (*Store, error) {
|
||||
db.Close()
|
||||
return nil, fmt.Errorf("resolve owner: %w", err)
|
||||
}
|
||||
return &Store{db: db, ownerID: ownerID}, nil
|
||||
return &Store{db: db, ownerID: ownerID, coverDir: coverDir}, nil
|
||||
}
|
||||
|
||||
// seedOwner makes sure the configured owner exists as exactly one readers row.
|
||||
@@ -550,40 +569,96 @@ func scanSeries(scan func(...any) error) (Series, error) {
|
||||
// Close releases the underlying database handle.
|
||||
func (s *Store) Close() error { return s.db.Close() }
|
||||
|
||||
// GetKaganeCover returns one persisted cover. Missing covers are reported with
|
||||
// ok=false rather than as an error so the web handler can fetch them once.
|
||||
func (s *Store) GetKaganeCover(imageID string) ([]byte, string, bool, error) {
|
||||
var (
|
||||
body []byte
|
||||
contentType string
|
||||
)
|
||||
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 kaganeCoverSourceURL(imageID string) string {
|
||||
return "https://kagane.to/api/v2/image/" + imageID + "/compressed"
|
||||
}
|
||||
|
||||
func (s *Store) getCover(sourceURL string) ([]byte, string, bool, error) {
|
||||
address := coverSourceAddress(sourceURL)
|
||||
var relativePath, contentType string
|
||||
err := s.db.QueryRow(
|
||||
`SELECT body, content_type FROM covers WHERE image_id = $1`, imageID,
|
||||
).Scan(&body, &contentType)
|
||||
`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 kagane cover %q: %w", imageID, err)
|
||||
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
|
||||
}
|
||||
|
||||
// PutKaganeCover persists one fetched cover. Image ids are immutable, so a
|
||||
// later fetch cannot replace the bytes already made durable.
|
||||
func (s *Store) PutKaganeCover(imageID string, body []byte, contentType string) error {
|
||||
func (s *Store) putCover(sourceURL string, body []byte, contentType string) error {
|
||||
if !IsKaganeCoverContentType(contentType) {
|
||||
return fmt.Errorf("put kagane cover %q: unsupported content type %q", imageID, contentType)
|
||||
return fmt.Errorf("put cover %q: unsupported content type %q", sourceURL, contentType)
|
||||
}
|
||||
address := coverSourceAddress(sourceURL)
|
||||
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 (image_id, body, content_type)
|
||||
INSERT INTO covers (address, path, content_type)
|
||||
VALUES ($1, $2, $3)
|
||||
ON CONFLICT (image_id) DO NOTHING`, imageID, body, contentType); err != nil {
|
||||
return fmt.Errorf("put kagane cover %q: %w", imageID, err)
|
||||
ON CONFLICT (address) DO NOTHING`, address, relativePath, contentType); err != nil {
|
||||
return fmt.Errorf("record cover %q: %w", address, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetKaganeCover returns one persisted cover. Missing covers are reported with
|
||||
// ok=false rather than as an error so the web handler can fetch them once.
|
||||
func (s *Store) GetKaganeCover(imageID string) ([]byte, string, bool, error) {
|
||||
return s.getCover(kaganeCoverSourceURL(imageID))
|
||||
}
|
||||
|
||||
// PutKaganeCover persists one fetched cover. The source URL's content address
|
||||
// makes each stored object immutable, so later writes for that URL are ignored.
|
||||
func (s *Store) PutKaganeCover(imageID string, body []byte, contentType string) error {
|
||||
return s.putCover(kaganeCoverSourceURL(imageID), body, contentType)
|
||||
}
|
||||
|
||||
// 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).
|
||||
|
||||
@@ -4,7 +4,9 @@ import (
|
||||
"bytes"
|
||||
"crypto/sha256"
|
||||
"database/sql"
|
||||
"encoding/hex"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
@@ -21,7 +23,7 @@ var testOwner = Owner{DiscordID: "test-owner", TokenHash: sha256.Sum256([]byte("
|
||||
|
||||
func newTestStore(t *testing.T) *Store {
|
||||
t.Helper()
|
||||
store, err := Open(pgtest.URL(t), testOwner)
|
||||
store, err := Open(pgtest.URL(t), testOwner, t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
@@ -46,7 +48,8 @@ func secondReader(t *testing.T, s *Store) int64 {
|
||||
// error, and must leave the rows alone.
|
||||
func TestOpenIsIdempotent(t *testing.T) {
|
||||
url := pgtest.URL(t)
|
||||
first, err := Open(url, testOwner)
|
||||
coverDir := t.TempDir()
|
||||
first, err := Open(url, testOwner, coverDir)
|
||||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
@@ -57,7 +60,7 @@ func TestOpenIsIdempotent(t *testing.T) {
|
||||
}
|
||||
first.Close()
|
||||
|
||||
second, err := Open(url, testOwner)
|
||||
second, err := Open(url, testOwner, coverDir)
|
||||
if err != nil {
|
||||
t.Fatalf("reopen: %v", err)
|
||||
}
|
||||
@@ -109,7 +112,8 @@ func TestReaderTokenInfo(t *testing.T) {
|
||||
// epoch-0 hash only while the row has never been rotated.
|
||||
func TestRotateTokenInvalidatesOldAndSurvivesRestart(t *testing.T) {
|
||||
url := pgtest.URL(t)
|
||||
store, err := Open(url, testOwner)
|
||||
coverDir := t.TempDir()
|
||||
store, err := Open(url, testOwner, coverDir)
|
||||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
@@ -141,7 +145,7 @@ func TestRotateTokenInvalidatesOldAndSurvivesRestart(t *testing.T) {
|
||||
}
|
||||
store.Close()
|
||||
|
||||
reopened, err := Open(url, testOwner)
|
||||
reopened, err := Open(url, testOwner, coverDir)
|
||||
if err != nil {
|
||||
t.Fatalf("reopen: %v", err)
|
||||
}
|
||||
@@ -668,7 +672,7 @@ func TestMigration0002BackfillsExistingBookmarks(t *testing.T) {
|
||||
// Bring it current through the production path: Open runs the schema to
|
||||
// 0003, seeds the owner, then applies 0004 which attaches this row. 0002
|
||||
// must have backfilled the series row, not lost data.
|
||||
st, err := Open(url, testOwner)
|
||||
st, err := Open(url, testOwner, t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatalf("Open after migrate: %v", err)
|
||||
}
|
||||
@@ -697,6 +701,48 @@ func TestMigration0002BackfillsExistingBookmarks(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigration0008DropsLegacyCoverRows(t *testing.T) {
|
||||
db, err := sql.Open("pgx", pgtest.URL(t))
|
||||
if err != nil {
|
||||
t.Fatalf("open: %v", err)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
if err := migrate(db, readersMigration); err != nil {
|
||||
t.Fatalf("migrate readers: %v", err)
|
||||
}
|
||||
if err := seedOwner(db, testOwner); err != nil {
|
||||
t.Fatalf("seed owner: %v", err)
|
||||
}
|
||||
if err := migrate(db, 7); err != nil {
|
||||
t.Fatalf("migrate legacy covers: %v", err)
|
||||
}
|
||||
if _, err := db.Exec(`
|
||||
INSERT INTO covers (image_id, body, content_type)
|
||||
VALUES ('legacy-image', 'legacy-bytes', 'image/jpeg')`); err != nil {
|
||||
t.Fatalf("seed legacy cover: %v", err)
|
||||
}
|
||||
if err := migrate(db, 0); err != nil {
|
||||
t.Fatalf("migrate filesystem covers: %v", err)
|
||||
}
|
||||
|
||||
var count int
|
||||
if err := db.QueryRow(`SELECT count(*) FROM covers`).Scan(&count); err != nil {
|
||||
t.Fatalf("count covers: %v", err)
|
||||
}
|
||||
if count != 0 {
|
||||
t.Fatalf("legacy covers = %d, want 0", count)
|
||||
}
|
||||
var bodyColumn int
|
||||
if err := db.QueryRow(`SELECT count(*) FROM information_schema.columns
|
||||
WHERE table_name = 'covers' AND column_name = 'body'`).Scan(&bodyColumn); err != nil {
|
||||
t.Fatalf("cover columns: %v", err)
|
||||
}
|
||||
if bodyColumn != 0 {
|
||||
t.Fatal("legacy covers table still has body column")
|
||||
}
|
||||
}
|
||||
|
||||
// readSeries reads the series row directly, for asserting on what Upsert
|
||||
// actually stored rather than what the joined Bookmark reports.
|
||||
func readSeries(t *testing.T, s *Store, site, seriesID string) Series {
|
||||
@@ -893,14 +939,15 @@ func TestDueForLatestCheckExcludesOrphanSeries(t *testing.T) {
|
||||
// rotations.
|
||||
func TestSeedOwnerIdempotentAndRefreshesTokenHash(t *testing.T) {
|
||||
url := pgtest.URL(t)
|
||||
first, err := Open(url, Owner{DiscordID: "owner", TokenHash: sha256.Sum256([]byte("hash-v1"))})
|
||||
coverDir := t.TempDir()
|
||||
first, err := Open(url, Owner{DiscordID: "owner", TokenHash: sha256.Sum256([]byte("hash-v1"))}, coverDir)
|
||||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
ownerID := first.OwnerID()
|
||||
first.Close()
|
||||
|
||||
second, err := Open(url, Owner{DiscordID: "owner", TokenHash: sha256.Sum256([]byte("hash-v2"))})
|
||||
second, err := Open(url, Owner{DiscordID: "owner", TokenHash: sha256.Sum256([]byte("hash-v2"))}, coverDir)
|
||||
if err != nil {
|
||||
t.Fatalf("reopen: %v", err)
|
||||
}
|
||||
@@ -957,7 +1004,7 @@ func TestMigration0004AttachesBookmarksToOwner(t *testing.T) {
|
||||
t.Fatalf("migrate to 0002: %v", err)
|
||||
}
|
||||
|
||||
st, err := Open(url, testOwner)
|
||||
st, err := Open(url, testOwner, t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
@@ -1233,7 +1280,8 @@ func TestTwoReadersShareOneSeriesWithIndependentProgress(t *testing.T) {
|
||||
}
|
||||
func TestKaganeCoverPersistsAcrossReopen(t *testing.T) {
|
||||
url := pgtest.URL(t)
|
||||
first, err := Open(url, testOwner)
|
||||
coverDir := t.TempDir()
|
||||
first, err := Open(url, testOwner, coverDir)
|
||||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
@@ -1245,7 +1293,7 @@ func TestKaganeCoverPersistsAcrossReopen(t *testing.T) {
|
||||
t.Fatalf("close first store: %v", err)
|
||||
}
|
||||
|
||||
second, err := Open(url, testOwner)
|
||||
second, err := Open(url, testOwner, coverDir)
|
||||
if err != nil {
|
||||
t.Fatalf("reopen: %v", err)
|
||||
}
|
||||
@@ -1258,3 +1306,49 @@ func TestKaganeCoverPersistsAcrossReopen(t *testing.T) {
|
||||
t.Fatalf("stored cover = (%q, %q, %v), want (%q, image/webp, true)", got, contentType, ok, body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenRequiresCoverDirectory(t *testing.T) {
|
||||
if _, err := Open(pgtest.URL(t), testOwner, ""); err == nil || !strings.Contains(err.Error(), "cover directory is required") {
|
||||
t.Fatalf("Open without cover directory = %v, want required-directory error", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestKaganeCoverIsContentAddressedOnFilesystem(t *testing.T) {
|
||||
url := pgtest.URL(t)
|
||||
coverDir := t.TempDir()
|
||||
first, err := Open(url, testOwner, coverDir)
|
||||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
body := []byte("stored-cover")
|
||||
const imageID = "019fe11a-84c3-7fc3-a84b-88787374b617"
|
||||
if err := first.PutKaganeCover(imageID, body, "image/webp"); err != nil {
|
||||
first.Close()
|
||||
t.Fatalf("PutKaganeCover: %v", err)
|
||||
}
|
||||
defer first.Close()
|
||||
|
||||
sourceURL := "https://kagane.to/api/v2/image/" + imageID + "/compressed"
|
||||
addressBytes := sha256.Sum256([]byte(sourceURL))
|
||||
address := hex.EncodeToString(addressBytes[:])
|
||||
wantPath := filepath.Join(address[:2], address[2:4], address)
|
||||
|
||||
var gotPath, contentType string
|
||||
if err := first.db.QueryRow(`SELECT path, content_type FROM covers WHERE address = $1`, address).Scan(&gotPath, &contentType); err != nil {
|
||||
t.Fatalf("cover row: %v", err)
|
||||
}
|
||||
if gotPath != wantPath || contentType != "image/webp" {
|
||||
t.Fatalf("cover row = (%q, %q), want (%q, image/webp)", gotPath, contentType, wantPath)
|
||||
}
|
||||
if got, err := os.ReadFile(filepath.Join(coverDir, gotPath)); err != nil || !bytes.Equal(got, body) {
|
||||
t.Fatalf("cover file = (%q, %v), want (%q, nil)", got, err, body)
|
||||
}
|
||||
|
||||
var bodyColumn int
|
||||
if err := first.db.QueryRow(`SELECT count(*) FROM information_schema.columns WHERE table_name = 'covers' AND column_name = 'body'`).Scan(&bodyColumn); err != nil {
|
||||
t.Fatalf("cover columns: %v", err)
|
||||
}
|
||||
if bodyColumn != 0 {
|
||||
t.Fatalf("covers still has body column")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user