// Package store is the read-only data access layer for the OpenGoods API. // It only issues SELECT queries; all writes happen in the Python ingestion path. package store import ( "context" "crypto/sha1" "encoding/hex" "encoding/json" "errors" "strconv" "strings" "time" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" "github.com/baicai2026-baicai/goods/api/internal/cache" ) // Cache TTLs for the public read cache. Product details are far less volatile // than search result sets, so they live longer; both are also invalidated // wholesale whenever ingestion bumps the cache epoch. const ( productCacheTTL = 24 * time.Hour searchCacheTTL = time.Hour ) // ErrNotFound is returned when a requested row does not exist. var ErrNotFound = errors.New("not found") // Store wraps a PostgreSQL connection pool and an optional read cache. type Store struct { pool *pgxpool.Pool cache *cache.Cache } // New constructs a Store from an existing pgx pool. func New(pool *pgxpool.Pool) *Store { return &Store{pool: pool} } // WithCache attaches a Redis-backed read cache. A nil or disabled cache leaves // the Store reading straight from PostgreSQL. func (s *Store) WithCache(c *cache.Cache) *Store { s.cache = c return s } // cacheGet reads a cached JSON value into dest, reporting a hit. It is a no-op // miss when no cache is attached. func (s *Store) cacheGet(ctx context.Context, suffix string, dest any) bool { if s.cache == nil { return false } return s.cache.GetJSON(ctx, suffix, dest) } // cacheSet stores a JSON value when a cache is attached. func (s *Store) cacheSet(ctx context.Context, suffix string, val any, ttl time.Duration) { if s.cache == nil { return } s.cache.SetJSON(ctx, suffix, val, ttl) } // Ping verifies database connectivity. func (s *Store) Ping(ctx context.Context) error { return s.pool.Ping(ctx) } // PublicStats summarizes the public catalog for the homepage. type PublicStats struct { Total int `json:"total"` Qualified int `json:"qualified"` MinScore float64 `json:"min_score"` } // Stats returns active-product totals and the number of qualified records whose // quality_score meets minScore. func (s *Store) Stats(ctx context.Context, minScore float64) (PublicStats, error) { st := PublicStats{MinScore: minScore} err := s.pool.QueryRow(ctx, ` SELECT count(*) FILTER (WHERE status = 'active'), count(*) FILTER (WHERE status = 'active' AND quality_score >= $1) FROM product`, minScore).Scan(&st.Total, &st.Qualified) return st, err } // Barcode is one GS1 trade item number attached to a product. type Barcode struct { GTIN string `json:"gtin"` GTINType string `json:"gtin_type"` PackLevel string `json:"pack_level"` Region *string `json:"region"` IsPrimary bool `json:"is_primary"` } // ProductSpec is one labeled spec line for a non-food product, rendered from // the product's attributes JSONB against its archive kind's field template. type ProductSpec struct { Key string `json:"key"` Label string `json:"label"` Value string `json:"value"` Unit string `json:"unit,omitempty"` } // Product is the full public view of a product. type Product struct { ID string `json:"id"` GTIN *string `json:"gtin"` Name string `json:"name"` Brand *string `json:"brand"` CategoryPath *string `json:"category_path"` GPCBrickCode *string `json:"gpc_brick_code"` ArchiveKind string `json:"archive_kind"` NetContentValue *float64 `json:"net_content_value"` NetContentUnit *string `json:"net_content_unit"` CountryOfOrigin *string `json:"country_of_origin"` QualityScore float64 `json:"quality_score"` Barcodes []Barcode `json:"barcodes"` Specs []ProductSpec `json:"specs,omitempty"` Nutriments map[string]any `json:"nutriments,omitempty"` NutritionBasis *string `json:"nutrition_basis,omitempty"` NutriScore *string `json:"nutri_score,omitempty"` Ingredients *string `json:"ingredients_text,omitempty"` Allergens []string `json:"allergens,omitempty"` Additives []string `json:"additives,omitempty"` MSRP []MSRP `json:"msrp,omitempty"` } // ProductBarcodes returns every barcode attached to a product, primary first. func (s *Store) ProductBarcodes(ctx context.Context, productID string) ([]Barcode, error) { rows, err := s.pool.Query(ctx, `SELECT gtin, gtin_type, pack_level, region, is_primary FROM product_barcode WHERE product_id = $1 ORDER BY is_primary DESC, gtin`, productID) if err != nil { return nil, err } defer rows.Close() out := []Barcode{} for rows.Next() { var b Barcode if err := rows.Scan(&b.GTIN, &b.GTINType, &b.PackLevel, &b.Region, &b.IsPrimary); err != nil { return nil, err } out = append(out, b) } return out, rows.Err() } // ProductSummary is a lightweight row used in search/listing responses. type ProductSummary struct { ID string `json:"id"` GTIN *string `json:"gtin"` Name string `json:"name"` Brand *string `json:"brand"` CategoryPath *string `json:"category_path"` Country *string `json:"country_of_origin"` QualityScore float64 `json:"quality_score"` Score *float64 `json:"score,omitempty"` } // SearchFilters bundles the optional filters accepted by SearchProducts. type SearchFilters struct { Query string // fuzzy name / barcode query Category string // ltree path; matches the subtree Brand string // fuzzy brand name Country string // country_of_origin prefix (case-insensitive) } const productSelect = ` SELECT p.id, p.gtin, p.name, b.name, c.path::text, p.gpc_brick_code, COALESCE(c.archive_kind, 'generic'), p.attributes, p.net_content_value, p.net_content_unit, p.country_of_origin, p.quality_score, f.nutriments, f.nutrition_basis, f.nutri_score, f.ingredients_text, f.allergens, f.additives FROM product p LEFT JOIN brand b ON b.id = p.brand_id LEFT JOIN category c ON c.id = p.category_id LEFT JOIN food_detail f ON f.product_id = p.id ` func scanProduct(row pgx.Row) (*Product, []byte, error) { var p Product var attributes []byte err := row.Scan( &p.ID, &p.GTIN, &p.Name, &p.Brand, &p.CategoryPath, &p.GPCBrickCode, &p.ArchiveKind, &attributes, &p.NetContentValue, &p.NetContentUnit, &p.CountryOfOrigin, &p.QualityScore, &p.Nutriments, &p.NutritionBasis, &p.NutriScore, &p.Ingredients, &p.Allergens, &p.Additives, ) if errors.Is(err, pgx.ErrNoRows) { return nil, nil, ErrNotFound } if err != nil { return nil, nil, err } return &p, attributes, nil } // buildSpecs renders the labeled, ordered spec list for a non-food product from // its attributes JSONB against its archive kind's field template. func (s *Store) buildSpecs(ctx context.Context, kind string, attributes []byte) ([]ProductSpec, error) { if kind == "" || kind == "food" || len(attributes) == 0 { return nil, nil } attrs := map[string]any{} if err := json.Unmarshal(attributes, &attrs); err != nil || len(attrs) == 0 { return nil, nil } rows, err := s.pool.Query(ctx, "SELECT field_key, label_zh, COALESCE(unit, '') FROM kind_field WHERE kind = $1 ORDER BY sort_order, field_key", kind) if err != nil { return nil, err } defer rows.Close() specs := []ProductSpec{} for rows.Next() { var key, label, unit string if err := rows.Scan(&key, &label, &unit); err != nil { return nil, err } v, ok := attrs[key] if !ok || v == nil { continue } val := stringifyAttr(v) if val == "" { continue } specs = append(specs, ProductSpec{Key: key, Label: label, Value: val, Unit: unit}) } return specs, rows.Err() } // stringifyAttr renders a JSON attribute value as display text. func stringifyAttr(v any) string { switch t := v.(type) { case string: return t case float64: return strconv.FormatFloat(t, 'f', -1, 64) case bool: if t { return "是" } return "否" case []any: parts := make([]string, 0, len(t)) for _, e := range t { parts = append(parts, stringifyAttr(e)) } return strings.Join(parts, "、") default: return "" } } // ProductByGTIN looks up an active product by any of its barcodes. func (s *Store) ProductByGTIN(ctx context.Context, gtin string) (*Product, error) { const suffix = "prod:gtin:" if cached := new(Product); s.cacheGet(ctx, suffix+gtin, cached) { return cached, nil } row := s.pool.QueryRow(ctx, productSelect+` WHERE p.status = 'active' AND (p.gtin = $1 OR EXISTS ( SELECT 1 FROM product_barcode pb WHERE pb.product_id = p.id AND pb.gtin = $1)) LIMIT 1`, gtin) p, attrs, err := scanProduct(row) if err != nil { return nil, err } if p.Barcodes, err = s.ProductBarcodes(ctx, p.ID); err != nil { return nil, err } if p.Specs, err = s.buildSpecs(ctx, p.ArchiveKind, attrs); err != nil { return nil, err } if p.MSRP, err = s.ListMSRP(ctx, p.ID); err != nil { return nil, err } s.cacheSet(ctx, suffix+gtin, p, productCacheTTL) return p, nil } // ProductByID looks up a product by its UUID. func (s *Store) ProductByID(ctx context.Context, id string) (*Product, error) { const suffix = "prod:id:" if cached := new(Product); s.cacheGet(ctx, suffix+id, cached) { return cached, nil } row := s.pool.QueryRow(ctx, productSelect+" WHERE p.id = $1", id) p, attrs, err := scanProduct(row) if err != nil { return nil, err } if p.Barcodes, err = s.ProductBarcodes(ctx, p.ID); err != nil { return nil, err } if p.Specs, err = s.buildSpecs(ctx, p.ArchiveKind, attrs); err != nil { return nil, err } if p.MSRP, err = s.ListMSRP(ctx, p.ID); err != nil { return nil, err } s.cacheSet(ctx, suffix+id, p, productCacheTTL) return p, nil } // fuzzyThreshold is the minimum word_similarity for a name to be considered a // fuzzy match. ~0.42 tolerates common typos (e.g. "choclate"→"Chocolate") // without returning unrelated products. const fuzzyThreshold = "0.42" // SearchProducts runs a trigram-fuzzy name search with optional category / // brand / country filters. When a query is present, matching is inclusive // (substring OR trigram-similar OR barcode), and results are ranked by name // similarity blended with quality_score so the best, most-complete records // surface first. Without a query, results are ordered by quality_score. func (s *Store) SearchProducts(ctx context.Context, f SearchFilters, limit, offset int) ([]ProductSummary, int, error) { suffix := searchCacheSuffix(f, limit, offset) if entry := new(searchCacheEntry); s.cacheGet(ctx, suffix, entry) { return entry.Items, entry.Total, nil } args := []any{} where := "WHERE p.status = 'active'" qIdx := 0 if f.Query != "" { args = append(args, f.Query) qIdx = len(args) q := "$" + strconv.Itoa(qIdx) where += ` AND (p.name ILIKE '%' || ` + q + ` || '%' OR word_similarity(` + q + `, p.name) >= ` + fuzzyThreshold + ` OR EXISTS (SELECT 1 FROM product_barcode pb WHERE pb.product_id = p.id AND pb.gtin ILIKE '%' || ` + q + ` || '%'))` } if f.Category != "" { args = append(args, f.Category) where += " AND c.path <@ $" + strconv.Itoa(len(args)) + "::ltree" } if f.Brand != "" { args = append(args, f.Brand) where += " AND b.name ILIKE '%' || $" + strconv.Itoa(len(args)) + " || '%'" } if f.Country != "" { args = append(args, f.Country) where += " AND p.country_of_origin ILIKE $" + strconv.Itoa(len(args)) + " || '%'" } from := `FROM product p LEFT JOIN brand b ON b.id = p.brand_id LEFT JOIN category c ON c.id = p.category_id ` var total int if err := s.pool.QueryRow(ctx, "SELECT count(*) "+from+where, args...).Scan(&total); err != nil { return nil, 0, err } // Ranking: when querying, similarity drives order, multiplied by a // quality factor floored at 0.5 so low-quality records aren't zeroed out. scoreExpr := "NULL::real" orderBy := "p.quality_score DESC, p.name" if f.Query != "" { q := "$" + strconv.Itoa(qIdx) scoreExpr = "word_similarity(" + q + ", p.name)" orderBy = scoreExpr + " * (0.5 + p.quality_score) DESC, p.quality_score DESC, p.name" } args = append(args, limit, offset) listSQL := "SELECT p.id, p.gtin, p.name, b.name, c.path::text, p.country_of_origin, p.quality_score, " + scoreExpr + " AS score " + from + where + " ORDER BY " + orderBy + " LIMIT $" + strconv.Itoa(len(args)-1) + " OFFSET $" + strconv.Itoa(len(args)) rows, err := s.pool.Query(ctx, listSQL, args...) if err != nil { return nil, 0, err } defer rows.Close() out := []ProductSummary{} for rows.Next() { var ps ProductSummary if err := rows.Scan(&ps.ID, &ps.GTIN, &ps.Name, &ps.Brand, &ps.CategoryPath, &ps.Country, &ps.QualityScore, &ps.Score); err != nil { return nil, 0, err } out = append(out, ps) } if err := rows.Err(); err != nil { return nil, 0, err } s.cacheSet(ctx, suffix, searchCacheEntry{Items: out, Total: total}, searchCacheTTL) return out, total, nil } // searchCacheEntry is the cached payload for a SearchProducts call. type searchCacheEntry struct { Items []ProductSummary `json:"items"` Total int `json:"total"` } // searchCacheSuffix derives a stable cache key from the full filter set and // paging window so distinct queries never collide. func searchCacheSuffix(f SearchFilters, limit, offset int) string { raw := strings.Join([]string{ f.Query, f.Category, f.Brand, f.Country, strconv.Itoa(limit), strconv.Itoa(offset), }, "\x1f") sum := sha1.Sum([]byte(raw)) return "search:" + hex.EncodeToString(sum[:]) } // Nutriments returns just the nutrition payload for a product. type Nutriments struct { ProductID string `json:"product_id"` Basis *string `json:"basis"` NutriScore *string `json:"nutri_score"` Values map[string]any `json:"values"` } // Nutriments fetches the nutrition facts of a product. func (s *Store) Nutriments(ctx context.Context, id string) (*Nutriments, error) { var n Nutriments n.ProductID = id err := s.pool.QueryRow(ctx, "SELECT nutriments, nutrition_basis, nutri_score FROM food_detail WHERE product_id = $1", id, ).Scan(&n.Values, &n.Basis, &n.NutriScore) if errors.Is(err, pgx.ErrNoRows) { return nil, ErrNotFound } if err != nil { return nil, err } return &n, nil } // MSRP is an official suggested retail price snapshot (never a purchase link). type MSRP struct { Amount float64 `json:"amount"` Currency string `json:"currency"` Region string `json:"region"` EffectiveDate *string `json:"effective_date"` SourceURL *string `json:"source_url"` Note *string `json:"note"` } // ListMSRP returns a product's suggested-retail-price snapshots, newest first. // Zero-amount entries are excluded as they carry no price information. func (s *Store) ListMSRP(ctx context.Context, id string) ([]MSRP, error) { rows, err := s.pool.Query(ctx, `SELECT amount, currency, region, effective_date::text, source_url, note FROM product_msrp WHERE product_id = $1 AND amount > 0 ORDER BY effective_date DESC NULLS LAST`, id) if err != nil { return nil, err } defer rows.Close() out := []MSRP{} for rows.Next() { var m MSRP if err := rows.Scan(&m.Amount, &m.Currency, &m.Region, &m.EffectiveDate, &m.SourceURL, &m.Note); err != nil { return nil, err } out = append(out, m) } return out, rows.Err() } // Brand is a public brand entry. type Brand struct { ID string `json:"id"` Name string `json:"name"` } // ListBrands returns brands ordered by name. func (s *Store) ListBrands(ctx context.Context, limit, offset int) ([]Brand, int, error) { var total int if err := s.pool.QueryRow(ctx, "SELECT count(*) FROM brand").Scan(&total); err != nil { return nil, 0, err } rows, err := s.pool.Query(ctx, "SELECT id, name FROM brand ORDER BY name LIMIT $1 OFFSET $2", limit, offset) if err != nil { return nil, 0, err } defer rows.Close() out := []Brand{} for rows.Next() { var b Brand if err := rows.Scan(&b.ID, &b.Name); err != nil { return nil, 0, err } out = append(out, b) } return out, total, rows.Err() } // Category is a node in the self-built category tree. type Category struct { ID string `json:"id"` NameZH string `json:"name_zh"` NameEN *string `json:"name_en"` Path string `json:"path"` GPCBrickCode *string `json:"gpc_brick_code"` Level int `json:"level"` ArchiveKind string `json:"archive_kind"` } // ListCategories returns the full category tree ordered by path. func (s *Store) ListCategories(ctx context.Context) ([]Category, error) { rows, err := s.pool.Query(ctx, "SELECT id, name_zh, name_en, path::text, gpc_brick_code, level, COALESCE(archive_kind, 'generic') FROM category ORDER BY path") if err != nil { return nil, err } defer rows.Close() out := []Category{} for rows.Next() { var c Category if err := rows.Scan(&c.ID, &c.NameZH, &c.NameEN, &c.Path, &c.GPCBrickCode, &c.Level, &c.ArchiveKind); err != nil { return nil, err } out = append(out, c) } return out, rows.Err() } // KindField describes one editable spec field for an archive kind. It drives // the dynamic contribution form (public, read-only view of the template). type KindField struct { Kind string `json:"kind"` FieldKey string `json:"field_key"` GroupLabel string `json:"group_label"` LabelZH string `json:"label_zh"` FieldType string `json:"field_type"` Unit *string `json:"unit"` Options []string `json:"options"` Placeholder *string `json:"placeholder"` SortOrder int `json:"sort_order"` Qualified bool `json:"qualified"` } // ListKindFields returns the ordered field template for one archive kind so the // public contribution form can render kind-specific inputs. func (s *Store) ListKindFields(ctx context.Context, kind string) ([]KindField, error) { rows, err := s.pool.Query(ctx, ` SELECT kind, field_key, group_label, label_zh, field_type, unit, options, placeholder, sort_order, qualified FROM kind_field WHERE kind = $1 ORDER BY sort_order, field_key`, kind) if err != nil { return nil, err } defer rows.Close() out := []KindField{} for rows.Next() { var f KindField if err := rows.Scan(&f.Kind, &f.FieldKey, &f.GroupLabel, &f.LabelZH, &f.FieldType, &f.Unit, &f.Options, &f.Placeholder, &f.SortOrder, &f.Qualified); err != nil { return nil, err } out = append(out, f) } return out, rows.Err() } // APIKey is the minimal metadata the public API needs to authorize a caller. type APIKey struct { ID string Name string RateLimitPerMin int QuotaTotal int64 } // APIKeyByHash returns the active (non-revoked) key matching a SHA-256 hash, // or ErrNotFound if no such active key exists. func (s *Store) APIKeyByHash(ctx context.Context, hash string) (*APIKey, error) { var k APIKey err := s.pool.QueryRow(ctx, `SELECT id, name, rate_limit_per_min, quota_total FROM api_key WHERE key_hash = $1 AND revoked_at IS NULL`, hash, ).Scan(&k.ID, &k.Name, &k.RateLimitPerMin, &k.QuotaTotal) if errors.Is(err, pgx.ErrNoRows) { return nil, ErrNotFound } if err != nil { return nil, err } return &k, nil } // Source describes a data source with its license and trust weight. type Source struct { ID string `json:"id"` Name string `json:"name"` Homepage *string `json:"homepage"` License *string `json:"license"` TrustWeight float64 `json:"trust_weight"` } // SourceByID fetches a single data source. func (s *Store) SourceByID(ctx context.Context, id string) (*Source, error) { var src Source err := s.pool.QueryRow(ctx, "SELECT id, name, homepage, license, trust_weight FROM source WHERE id = $1", id, ).Scan(&src.ID, &src.Name, &src.Homepage, &src.License, &src.TrustWeight) if errors.Is(err, pgx.ErrNoRows) { return nil, ErrNotFound } if err != nil { return nil, err } return &src, nil }