// 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" "encoding/json" "errors" "strconv" "strings" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" ) // ErrNotFound is returned when a requested row does not exist. var ErrNotFound = errors.New("not found") // Store wraps a PostgreSQL connection pool. type Store struct { pool *pgxpool.Pool } // New constructs a Store from an existing pgx pool. func New(pool *pgxpool.Pool) *Store { return &Store{pool: pool} } // 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"` } // 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) { 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 } return p, nil } // ProductByID looks up a product by its UUID. func (s *Store) ProductByID(ctx context.Context, id string) (*Product, error) { 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 } 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) { 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) } return out, total, rows.Err() } // 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 all MSRP snapshots for a product. 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 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"` } // 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 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); err != nil { return nil, err } out = append(out, c) } 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 }