All checks were successful
Release / Scraper / Test (push) Successful in 10s
Release / UI / Build (push) Successful in 26s
Release / v2 / Build ui-v2 (push) Successful in 17s
Release / Scraper / Docker (push) Successful in 47s
Release / UI / Docker (push) Successful in 56s
CI / Scraper / Lint (pull_request) Successful in 7s
CI / Scraper / Test (pull_request) Successful in 8s
CI / UI / Build (pull_request) Successful in 16s
CI / Scraper / Docker Push (pull_request) Has been skipped
CI / UI / Docker Push (pull_request) Has been skipped
Release / v2 / Docker / ui-v2 (push) Successful in 56s
Release / v2 / Test backend (push) Successful in 4m35s
iOS CI / Build (pull_request) Successful in 4m28s
Release / v2 / Docker / backend (push) Successful in 1m29s
Release / v2 / Docker / runner (push) Successful in 1m39s
iOS CI / Test (pull_request) Successful in 9m51s
- backend/: Go API server and runner binaries with PocketBase + MinIO storage - ui-v2/: SvelteKit frontend rewrite - docker-compose-new.yml: compose file for the v2 stack - .gitea/workflows/release-v2.yaml: CI/CD for backend, runner, and ui-v2 Docker Hub images - scripts/pb-init.sh: migrate from wget to curl, add superuser bootstrap for fresh installs - .env.example: document DOCKER_BUILDKIT=1 for Colima users
269 lines
7.3 KiB
Go
269 lines
7.3 KiB
Go
// Package storage provides the concrete implementations of all bookstore and
|
|
// taskqueue interfaces backed by PocketBase (structured data) and MinIO (blobs).
|
|
//
|
|
// Entry point: NewStore(ctx, cfg, log) returns a *Store that satisfies every
|
|
// interface defined in bookstore and taskqueue.
|
|
package storage
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"net/http"
|
|
"net/url"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/libnovel/backend/internal/config"
|
|
"github.com/libnovel/backend/internal/domain"
|
|
)
|
|
|
|
// ErrNotFound is returned by single-record lookups when no record exists.
|
|
var ErrNotFound = errors.New("storage: record not found")
|
|
|
|
// pbClient is the internal PocketBase REST admin client.
|
|
type pbClient struct {
|
|
baseURL string
|
|
email string
|
|
password string
|
|
log *slog.Logger
|
|
|
|
mu sync.Mutex
|
|
token string
|
|
exp time.Time
|
|
}
|
|
|
|
func newPBClient(cfg config.PocketBase, log *slog.Logger) *pbClient {
|
|
return &pbClient{
|
|
baseURL: strings.TrimRight(cfg.URL, "/"),
|
|
email: cfg.AdminEmail,
|
|
password: cfg.AdminPassword,
|
|
log: log,
|
|
}
|
|
}
|
|
|
|
// authToken returns a valid admin auth token, refreshing it when expired.
|
|
func (c *pbClient) authToken(ctx context.Context) (string, error) {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
if c.token != "" && time.Now().Before(c.exp) {
|
|
return c.token, nil
|
|
}
|
|
|
|
body, _ := json.Marshal(map[string]string{
|
|
"identity": c.email,
|
|
"password": c.password,
|
|
})
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
|
|
c.baseURL+"/api/collections/_superusers/auth-with-password", bytes.NewReader(body))
|
|
if err != nil {
|
|
return "", fmt.Errorf("pb auth: build request: %w", err)
|
|
}
|
|
req.Header.Set("Content-Type", "application/json")
|
|
|
|
resp, err := http.DefaultClient.Do(req)
|
|
if err != nil {
|
|
return "", fmt.Errorf("pb auth: %w", err)
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
if resp.StatusCode != http.StatusOK {
|
|
raw, _ := io.ReadAll(resp.Body)
|
|
return "", fmt.Errorf("pb auth: status %d: %s", resp.StatusCode, string(raw))
|
|
}
|
|
|
|
var payload struct {
|
|
Token string `json:"token"`
|
|
}
|
|
if err := json.NewDecoder(resp.Body).Decode(&payload); err != nil {
|
|
return "", fmt.Errorf("pb auth: decode: %w", err)
|
|
}
|
|
c.token = payload.Token
|
|
c.exp = time.Now().Add(30 * time.Minute)
|
|
return c.token, nil
|
|
}
|
|
|
|
// do executes an authenticated PocketBase REST request.
|
|
func (c *pbClient) do(ctx context.Context, method, path string, body io.Reader) (*http.Response, error) {
|
|
tok, err := c.authToken(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, body)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("pb: build request %s %s: %w", method, path, err)
|
|
}
|
|
req.Header.Set("Authorization", tok)
|
|
if body != nil {
|
|
req.Header.Set("Content-Type", "application/json")
|
|
}
|
|
|
|
resp, err := http.DefaultClient.Do(req)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("pb: %s %s: %w", method, path, err)
|
|
}
|
|
return resp, nil
|
|
}
|
|
|
|
// get is a convenience wrapper that decodes a JSON response into v.
|
|
func (c *pbClient) get(ctx context.Context, path string, v any) error {
|
|
resp, err := c.do(ctx, http.MethodGet, path, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode == http.StatusNotFound {
|
|
return ErrNotFound
|
|
}
|
|
if resp.StatusCode >= 400 {
|
|
raw, _ := io.ReadAll(resp.Body)
|
|
return fmt.Errorf("pb GET %s: status %d: %s", path, resp.StatusCode, string(raw))
|
|
}
|
|
return json.NewDecoder(resp.Body).Decode(v)
|
|
}
|
|
|
|
// post creates a record and decodes the created record into v.
|
|
func (c *pbClient) post(ctx context.Context, path string, payload, v any) error {
|
|
b, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return fmt.Errorf("pb: marshal: %w", err)
|
|
}
|
|
resp, err := c.do(ctx, http.MethodPost, path, bytes.NewReader(b))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode >= 400 {
|
|
raw, _ := io.ReadAll(resp.Body)
|
|
return fmt.Errorf("pb POST %s: status %d: %s", path, resp.StatusCode, string(raw))
|
|
}
|
|
if v != nil {
|
|
return json.NewDecoder(resp.Body).Decode(v)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// patch updates a record.
|
|
func (c *pbClient) patch(ctx context.Context, path string, payload any) error {
|
|
b, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return fmt.Errorf("pb: marshal: %w", err)
|
|
}
|
|
resp, err := c.do(ctx, http.MethodPatch, path, bytes.NewReader(b))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode >= 400 {
|
|
raw, _ := io.ReadAll(resp.Body)
|
|
return fmt.Errorf("pb PATCH %s: status %d: %s", path, resp.StatusCode, string(raw))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// delete removes a record.
|
|
func (c *pbClient) delete(ctx context.Context, path string) error {
|
|
resp, err := c.do(ctx, http.MethodDelete, path, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode == http.StatusNotFound {
|
|
return ErrNotFound
|
|
}
|
|
if resp.StatusCode >= 400 {
|
|
raw, _ := io.ReadAll(resp.Body)
|
|
return fmt.Errorf("pb DELETE %s: status %d: %s", path, resp.StatusCode, string(raw))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// listAll fetches all pages of a collection. PocketBase returns at most 200
|
|
// records per page; we paginate until empty.
|
|
func (c *pbClient) listAll(ctx context.Context, collection string, filter, sort string) ([]json.RawMessage, error) {
|
|
var all []json.RawMessage
|
|
page := 1
|
|
for {
|
|
q := url.Values{
|
|
"page": {fmt.Sprintf("%d", page)},
|
|
"perPage": {"200"},
|
|
}
|
|
if filter != "" {
|
|
q.Set("filter", filter)
|
|
}
|
|
if sort != "" {
|
|
q.Set("sort", sort)
|
|
}
|
|
path := fmt.Sprintf("/api/collections/%s/records?%s", collection, q.Encode())
|
|
|
|
var result struct {
|
|
Items []json.RawMessage `json:"items"`
|
|
Page int `json:"page"`
|
|
Pages int `json:"totalPages"`
|
|
}
|
|
if err := c.get(ctx, path, &result); err != nil {
|
|
return nil, err
|
|
}
|
|
all = append(all, result.Items...)
|
|
if result.Page >= result.Pages {
|
|
break
|
|
}
|
|
page++
|
|
}
|
|
return all, nil
|
|
}
|
|
|
|
// claimRecord atomically claims the first pending record matching collection.
|
|
// It fetches the oldest pending record (filter + sort), then PATCHes it with
|
|
// the claim payload. Returns (nil, nil) when the queue is empty.
|
|
func (c *pbClient) claimRecord(ctx context.Context, collection, workerID string, extraClaim map[string]any) (json.RawMessage, error) {
|
|
q := url.Values{}
|
|
q.Set("filter", `status="pending"`)
|
|
q.Set("sort", "+started")
|
|
q.Set("perPage", "1")
|
|
path := fmt.Sprintf("/api/collections/%s/records?%s", collection, q.Encode())
|
|
|
|
var result struct {
|
|
Items []json.RawMessage `json:"items"`
|
|
}
|
|
if err := c.get(ctx, path, &result); err != nil {
|
|
return nil, fmt.Errorf("claimRecord list: %w", err)
|
|
}
|
|
if len(result.Items) == 0 {
|
|
return nil, nil // queue empty
|
|
}
|
|
|
|
var rec struct {
|
|
ID string `json:"id"`
|
|
}
|
|
if err := json.Unmarshal(result.Items[0], &rec); err != nil {
|
|
return nil, fmt.Errorf("claimRecord parse id: %w", err)
|
|
}
|
|
|
|
claim := map[string]any{
|
|
"status": string(domain.TaskStatusRunning),
|
|
"worker_id": workerID,
|
|
}
|
|
for k, v := range extraClaim {
|
|
claim[k] = v
|
|
}
|
|
|
|
claimPath := fmt.Sprintf("/api/collections/%s/records/%s", collection, rec.ID)
|
|
if err := c.patch(ctx, claimPath, claim); err != nil {
|
|
return nil, fmt.Errorf("claimRecord patch: %w", err)
|
|
}
|
|
|
|
// Re-fetch the updated record so caller has current state.
|
|
var updated json.RawMessage
|
|
if err := c.get(ctx, claimPath, &updated); err != nil {
|
|
return nil, fmt.Errorf("claimRecord re-fetch: %w", err)
|
|
}
|
|
return updated, nil
|
|
}
|