// 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 }