Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 5b9338175e | |||
| 98b5f1575d | |||
| 1d5f775d33 |
@@ -86,6 +86,20 @@ export const api = {
|
||||
request<{ status: string }>(`/products/${id}/msrp/${msrpId}`, {
|
||||
method: "DELETE",
|
||||
}),
|
||||
addBarcode: (id: string, body: unknown) =>
|
||||
request<import("./types").Barcode>(`/products/${id}/barcodes`, {
|
||||
method: "POST",
|
||||
body: JSON.stringify(body),
|
||||
}),
|
||||
deleteBarcode: (id: string, barcodeId: string) =>
|
||||
request<{ status: string }>(`/products/${id}/barcodes/${barcodeId}`, {
|
||||
method: "DELETE",
|
||||
}),
|
||||
setPrimaryBarcode: (id: string, barcodeId: string) =>
|
||||
request<import("./types").Barcode>(
|
||||
`/products/${id}/barcodes/${barcodeId}/primary`,
|
||||
{ method: "POST" },
|
||||
),
|
||||
listBrands: () =>
|
||||
request<{ items: import("./types").Brand[] }>("/brands"),
|
||||
listCategories: () =>
|
||||
|
||||
@@ -14,8 +14,16 @@ import {
|
||||
Trash2,
|
||||
AlertCircle,
|
||||
History,
|
||||
Star,
|
||||
} from "lucide-react";
|
||||
|
||||
const GTIN_TYPES = ["EAN13", "EAN8", "UPC", "ITF14", "GTIN14"];
|
||||
const PACK_LEVELS: { value: string; label: string }[] = [
|
||||
{ value: "each", label: "消费单元" },
|
||||
{ value: "case", label: "箱" },
|
||||
{ value: "pallet", label: "托盘" },
|
||||
];
|
||||
|
||||
const NUTRIMENT_KEYS: { key: string; label: string }[] = [
|
||||
{ key: "energy_kcal", label: "能量 (kcal)" },
|
||||
{ key: "energy_kj", label: "能量 (kJ)" },
|
||||
@@ -39,6 +47,9 @@ const ACTION_LABEL: Record<string, string> = {
|
||||
delete_image: "删除图片",
|
||||
add_msrp: "新增建议零售价",
|
||||
delete_msrp: "删除建议零售价",
|
||||
add_barcode: "新增条码",
|
||||
delete_barcode: "删除条码",
|
||||
set_primary_barcode: "设为主条码",
|
||||
};
|
||||
|
||||
function Card({
|
||||
@@ -400,6 +411,7 @@ export default function ProductDetail({
|
||||
</div>
|
||||
</Card>
|
||||
|
||||
<BarcodesCard product={d} onChange={reload} onError={setError} />
|
||||
<ImagesCard
|
||||
product={d}
|
||||
onChange={reload}
|
||||
@@ -435,6 +447,158 @@ export default function ProductDetail({
|
||||
);
|
||||
}
|
||||
|
||||
function BarcodesCard({
|
||||
product,
|
||||
onChange,
|
||||
onError,
|
||||
}: {
|
||||
product: Detail;
|
||||
onChange: () => void;
|
||||
onError: (m: string) => void;
|
||||
}) {
|
||||
const [gtin, setGtin] = useState("");
|
||||
const [gtinType, setGtinType] = useState("EAN13");
|
||||
const [packLevel, setPackLevel] = useState("each");
|
||||
const [region, setRegion] = useState("");
|
||||
const [busy, setBusy] = useState(false);
|
||||
|
||||
async function add() {
|
||||
if (!gtin.trim()) return;
|
||||
setBusy(true);
|
||||
try {
|
||||
await api.addBarcode(product.id, {
|
||||
gtin: gtin.trim(),
|
||||
gtin_type: gtinType,
|
||||
pack_level: packLevel,
|
||||
region: region.trim() || null,
|
||||
is_primary: false,
|
||||
});
|
||||
setGtin("");
|
||||
setRegion("");
|
||||
onChange();
|
||||
} catch (e) {
|
||||
onError(e instanceof Error ? e.message : "添加失败");
|
||||
} finally {
|
||||
setBusy(false);
|
||||
}
|
||||
}
|
||||
async function remove(barcodeId: string) {
|
||||
try {
|
||||
await api.deleteBarcode(product.id, barcodeId);
|
||||
onChange();
|
||||
} catch (e) {
|
||||
onError(e instanceof Error ? e.message : "删除失败");
|
||||
}
|
||||
}
|
||||
async function makePrimary(barcodeId: string) {
|
||||
try {
|
||||
await api.setPrimaryBarcode(product.id, barcodeId);
|
||||
onChange();
|
||||
} catch (e) {
|
||||
onError(e instanceof Error ? e.message : "设置失败");
|
||||
}
|
||||
}
|
||||
|
||||
return (
|
||||
<Card title="条码(一品多码,主条码镜像到 GTIN)">
|
||||
<div className="mb-3 space-y-2">
|
||||
{product.barcodes.length === 0 && (
|
||||
<span className="text-sm text-gray-400">暂无条码</span>
|
||||
)}
|
||||
{product.barcodes.map((b) => (
|
||||
<div
|
||||
key={b.id}
|
||||
className="flex items-center gap-3 rounded border border-gray-100 bg-gray-50 px-3 py-2 text-sm"
|
||||
>
|
||||
<button
|
||||
onClick={() => !b.is_primary && makePrimary(b.id)}
|
||||
title={b.is_primary ? "主条码" : "设为主条码"}
|
||||
disabled={b.is_primary}
|
||||
className={
|
||||
b.is_primary
|
||||
? "text-amber-500"
|
||||
: "text-gray-300 hover:text-amber-500"
|
||||
}
|
||||
>
|
||||
<Star
|
||||
className="h-4 w-4"
|
||||
fill={b.is_primary ? "currentColor" : "none"}
|
||||
/>
|
||||
</button>
|
||||
<span className="font-mono font-medium text-gray-800">
|
||||
{b.gtin}
|
||||
</span>
|
||||
<span className="rounded bg-gray-200 px-1.5 py-0.5 text-[11px] text-gray-600">
|
||||
{b.gtin_type}
|
||||
</span>
|
||||
<span className="text-gray-500">
|
||||
{PACK_LEVELS.find((p) => p.value === b.pack_level)?.label ||
|
||||
b.pack_level}
|
||||
</span>
|
||||
<span className="flex-1 text-gray-400">{b.region || ""}</span>
|
||||
<button
|
||||
onClick={() => remove(b.id)}
|
||||
className="text-gray-400 hover:text-red-600"
|
||||
>
|
||||
<Trash2 className="h-4 w-4" />
|
||||
</button>
|
||||
</div>
|
||||
))}
|
||||
</div>
|
||||
<div className="flex flex-wrap items-end gap-2">
|
||||
<Field label="条码 (GTIN)">
|
||||
<input
|
||||
className="w-44 rounded border border-gray-300 px-3 py-2 text-sm"
|
||||
value={gtin}
|
||||
onChange={(e) => setGtin(e.target.value)}
|
||||
placeholder="8/12/13/14 位"
|
||||
/>
|
||||
</Field>
|
||||
<Field label="类型">
|
||||
<select
|
||||
className="rounded border border-gray-300 px-2 py-2 text-sm"
|
||||
value={gtinType}
|
||||
onChange={(e) => setGtinType(e.target.value)}
|
||||
>
|
||||
{GTIN_TYPES.map((t) => (
|
||||
<option key={t} value={t}>
|
||||
{t}
|
||||
</option>
|
||||
))}
|
||||
</select>
|
||||
</Field>
|
||||
<Field label="包装层级">
|
||||
<select
|
||||
className="rounded border border-gray-300 px-2 py-2 text-sm"
|
||||
value={packLevel}
|
||||
onChange={(e) => setPackLevel(e.target.value)}
|
||||
>
|
||||
{PACK_LEVELS.map((p) => (
|
||||
<option key={p.value} value={p.value}>
|
||||
{p.label}
|
||||
</option>
|
||||
))}
|
||||
</select>
|
||||
</Field>
|
||||
<Field label="地区(可选)">
|
||||
<input
|
||||
className="w-20 rounded border border-gray-300 px-3 py-2 text-sm"
|
||||
value={region}
|
||||
onChange={(e) => setRegion(e.target.value.toUpperCase())}
|
||||
/>
|
||||
</Field>
|
||||
<button
|
||||
onClick={add}
|
||||
disabled={busy}
|
||||
className="flex items-center gap-1 rounded bg-gray-700 px-3 py-2 text-sm text-white hover:bg-gray-800 disabled:opacity-60"
|
||||
>
|
||||
<Plus className="h-4 w-4" /> 添加
|
||||
</button>
|
||||
</div>
|
||||
</Card>
|
||||
);
|
||||
}
|
||||
|
||||
function ImagesCard({
|
||||
product,
|
||||
onChange,
|
||||
|
||||
@@ -10,6 +10,15 @@ export interface ProductRow {
|
||||
updated_at: string;
|
||||
}
|
||||
|
||||
export interface Barcode {
|
||||
id: string;
|
||||
gtin: string;
|
||||
gtin_type: string;
|
||||
pack_level: string;
|
||||
region: string | null;
|
||||
is_primary: boolean;
|
||||
}
|
||||
|
||||
export interface ProductImage {
|
||||
id: string;
|
||||
url: string;
|
||||
@@ -47,6 +56,7 @@ export interface ProductDetail {
|
||||
nutrition_basis: string | null;
|
||||
serving_size: string | null;
|
||||
nutri_score: string | null;
|
||||
barcodes: Barcode[];
|
||||
images: ProductImage[];
|
||||
msrp: MSRP[];
|
||||
missing: string[];
|
||||
|
||||
@@ -15,6 +15,7 @@ import (
|
||||
|
||||
"github.com/baicai2026-baicai/goods/api/internal/adminstore"
|
||||
"github.com/baicai2026-baicai/goods/api/internal/auth"
|
||||
"github.com/baicai2026-baicai/goods/api/internal/gtin"
|
||||
)
|
||||
|
||||
// Handler holds the admin dependencies.
|
||||
@@ -62,6 +63,9 @@ func (h *Handler) Router() http.Handler {
|
||||
r.Delete("/api/products/{id}/images/{imageID}", h.DeleteImage)
|
||||
r.Post("/api/products/{id}/msrp", h.AddMSRP)
|
||||
r.Delete("/api/products/{id}/msrp/{msrpID}", h.DeleteMSRP)
|
||||
r.Post("/api/products/{id}/barcodes", h.AddBarcode)
|
||||
r.Delete("/api/products/{id}/barcodes/{barcodeID}", h.DeleteBarcode)
|
||||
r.Post("/api/products/{id}/barcodes/{barcodeID}/primary", h.SetPrimaryBarcode)
|
||||
r.Get("/api/brands", h.ListBrands)
|
||||
r.Get("/api/categories", h.ListCategories)
|
||||
|
||||
@@ -231,6 +235,68 @@ func (h *Handler) DeleteMSRP(w http.ResponseWriter, r *http.Request) {
|
||||
writeJSON(w, http.StatusOK, map[string]string{"status": "deleted"})
|
||||
}
|
||||
|
||||
// ---------- barcodes ----------
|
||||
|
||||
// AddBarcode validates and attaches a barcode to a product. A code already
|
||||
// owned by another product yields 409 with the conflicting product so the
|
||||
// operator can de-duplicate; an invalid GTIN yields 400.
|
||||
func (h *Handler) AddBarcode(w http.ResponseWriter, r *http.Request) {
|
||||
var in adminstore.BarcodeInput
|
||||
if err := json.NewDecoder(r.Body).Decode(&in); err != nil {
|
||||
writeError(w, http.StatusBadRequest, "bad_request", "invalid body")
|
||||
return
|
||||
}
|
||||
b, err := h.store.AddBarcode(r.Context(), chi.URLParam(r, "id"), auth.UserFrom(r.Context()), in)
|
||||
if h.handleBarcodeErr(w, err) {
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, b)
|
||||
}
|
||||
|
||||
// DeleteBarcode removes a barcode; a primary one is replaced automatically.
|
||||
func (h *Handler) DeleteBarcode(w http.ResponseWriter, r *http.Request) {
|
||||
err := h.store.DeleteBarcode(r.Context(), chi.URLParam(r, "id"), chi.URLParam(r, "barcodeID"), auth.UserFrom(r.Context()))
|
||||
if h.handleErr(w, err) {
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]string{"status": "deleted"})
|
||||
}
|
||||
|
||||
// SetPrimaryBarcode marks one barcode primary and mirrors it to product.gtin.
|
||||
func (h *Handler) SetPrimaryBarcode(w http.ResponseWriter, r *http.Request) {
|
||||
b, err := h.store.SetPrimaryBarcode(r.Context(), chi.URLParam(r, "id"), chi.URLParam(r, "barcodeID"), auth.UserFrom(r.Context()))
|
||||
if h.handleErr(w, err) {
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, b)
|
||||
}
|
||||
|
||||
// handleBarcodeErr maps barcode-specific errors (GTIN validation, ownership
|
||||
// conflict) to client-facing statuses, falling back to handleErr otherwise.
|
||||
func (h *Handler) handleBarcodeErr(w http.ResponseWriter, err error) bool {
|
||||
if err == nil {
|
||||
return false
|
||||
}
|
||||
var conflict *adminstore.ConflictError
|
||||
if errors.As(err, &conflict) {
|
||||
writeJSON(w, http.StatusConflict, map[string]any{
|
||||
"error": map[string]string{"code": "barcode_conflict", "message": err.Error()},
|
||||
"conflict": map[string]string{
|
||||
"gtin": conflict.GTIN,
|
||||
"product_id": conflict.ProductID,
|
||||
"product_name": conflict.ProductName,
|
||||
},
|
||||
})
|
||||
return true
|
||||
}
|
||||
if errors.Is(err, gtin.ErrEmpty) || errors.Is(err, gtin.ErrFormat) ||
|
||||
errors.Is(err, gtin.ErrCheck) || errors.Is(err, gtin.ErrRestricted) {
|
||||
writeError(w, http.StatusBadRequest, "invalid_gtin", err.Error())
|
||||
return true
|
||||
}
|
||||
return h.handleErr(w, err)
|
||||
}
|
||||
|
||||
// ListBrands returns brand options.
|
||||
func (h *Handler) ListBrands(w http.ResponseWriter, r *http.Request) {
|
||||
items, err := h.store.ListBrands(r.Context())
|
||||
|
||||
@@ -162,6 +162,7 @@ type ProductDetail struct {
|
||||
NutritionBasis *string `json:"nutrition_basis"`
|
||||
ServingSize *string `json:"serving_size"`
|
||||
NutriScore *string `json:"nutri_score"`
|
||||
Barcodes []Barcode `json:"barcodes"`
|
||||
Images []ProductImage `json:"images"`
|
||||
MSRP []MSRP `json:"msrp"`
|
||||
Missing []string `json:"missing"`
|
||||
@@ -207,6 +208,12 @@ WHERE p.id = $1`, id).Scan(
|
||||
d.Additives = []string{}
|
||||
}
|
||||
|
||||
bcs, err := s.listBarcodes(ctx, id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
d.Barcodes = bcs
|
||||
|
||||
imgs, err := s.listImages(ctx, id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
||||
@@ -0,0 +1,270 @@
|
||||
package adminstore
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
|
||||
"github.com/baicai2026-baicai/goods/api/internal/gtin"
|
||||
)
|
||||
|
||||
// Barcode is one GS1 trade item number attached to a product.
|
||||
type Barcode struct {
|
||||
ID string `json:"id"`
|
||||
GTIN string `json:"gtin"`
|
||||
GTINType string `json:"gtin_type"`
|
||||
PackLevel string `json:"pack_level"`
|
||||
Region *string `json:"region"`
|
||||
IsPrimary bool `json:"is_primary"`
|
||||
}
|
||||
|
||||
// BarcodeInput is the payload for attaching a barcode to a product.
|
||||
type BarcodeInput struct {
|
||||
GTIN string `json:"gtin"`
|
||||
GTINType string `json:"gtin_type"`
|
||||
PackLevel string `json:"pack_level"`
|
||||
Region *string `json:"region"`
|
||||
IsPrimary bool `json:"is_primary"`
|
||||
}
|
||||
|
||||
// ConflictError signals that a barcode is already attached to another product,
|
||||
// so the operator must de-duplicate instead of creating a clash.
|
||||
type ConflictError struct {
|
||||
GTIN string
|
||||
ProductID string
|
||||
ProductName string
|
||||
}
|
||||
|
||||
func (e *ConflictError) Error() string { return "条码已被其他商品占用:" + e.GTIN }
|
||||
|
||||
func validPackLevel(p string) string {
|
||||
switch p {
|
||||
case "each", "case", "pallet":
|
||||
return p
|
||||
default:
|
||||
return "each"
|
||||
}
|
||||
}
|
||||
|
||||
func validGTINType(t, normalized string) string {
|
||||
switch t {
|
||||
case "EAN8", "UPC", "EAN13", "ITF14", "GTIN14":
|
||||
return t
|
||||
default:
|
||||
return gtin.InferType(normalized)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Store) listBarcodes(ctx context.Context, productID string) ([]Barcode, error) {
|
||||
rows, err := s.pool.Query(ctx,
|
||||
`SELECT id, 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.ID, &b.GTIN, &b.GTINType, &b.PackLevel, &b.Region, &b.IsPrimary); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out = append(out, b)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// barcodeOwner returns the product currently owning a barcode, if any.
|
||||
func barcodeOwner(ctx context.Context, q pgx.Tx, code string) (productID, productName string, found bool, err error) {
|
||||
err = q.QueryRow(ctx,
|
||||
`SELECT pb.product_id, p.name FROM product_barcode pb
|
||||
JOIN product p ON p.id = pb.product_id WHERE pb.gtin = $1`, code).
|
||||
Scan(&productID, &productName)
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return "", "", false, nil
|
||||
}
|
||||
if err != nil {
|
||||
return "", "", false, err
|
||||
}
|
||||
return productID, productName, true, nil
|
||||
}
|
||||
|
||||
// AddBarcode validates and attaches a barcode to a product, recording audit.
|
||||
// A barcode already owned by another product yields a *ConflictError.
|
||||
func (s *Store) AddBarcode(ctx context.Context, productID, actor string, in BarcodeInput) (*Barcode, error) {
|
||||
code, err := gtin.Normalize(in.GTIN)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
tx, err := s.pool.Begin(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer tx.Rollback(ctx)
|
||||
|
||||
// Product must exist.
|
||||
var exists bool
|
||||
if err := tx.QueryRow(ctx, "SELECT EXISTS(SELECT 1 FROM product WHERE id=$1)", productID).Scan(&exists); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !exists {
|
||||
return nil, ErrNotFound
|
||||
}
|
||||
|
||||
// Globally unique: a barcode owned by any product (this one included)
|
||||
// is a conflict the operator must resolve by de-duplicating.
|
||||
if owner, name, found, err := barcodeOwner(ctx, tx, code); err != nil {
|
||||
return nil, err
|
||||
} else if found {
|
||||
return nil, &ConflictError{GTIN: code, ProductID: owner, ProductName: name}
|
||||
}
|
||||
|
||||
// Make this the primary barcode when requested or when none exists yet.
|
||||
makePrimary := in.IsPrimary
|
||||
if !makePrimary {
|
||||
var hasPrimary bool
|
||||
if err := tx.QueryRow(ctx,
|
||||
"SELECT EXISTS(SELECT 1 FROM product_barcode WHERE product_id=$1 AND is_primary)", productID).
|
||||
Scan(&hasPrimary); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
makePrimary = !hasPrimary
|
||||
}
|
||||
if makePrimary {
|
||||
if _, err := tx.Exec(ctx,
|
||||
"UPDATE product_barcode SET is_primary=false WHERE product_id=$1 AND is_primary", productID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
srcID, _ := s.manualSourceID(ctx, tx)
|
||||
var srcArg any
|
||||
if srcID != "" {
|
||||
srcArg = srcID
|
||||
}
|
||||
|
||||
var b Barcode
|
||||
err = tx.QueryRow(ctx, `
|
||||
INSERT INTO product_barcode (product_id, gtin, gtin_type, pack_level, region, is_primary, source_id)
|
||||
VALUES ($1,$2,$3,$4,$5,$6,$7)
|
||||
RETURNING id, gtin, gtin_type, pack_level, region, is_primary`,
|
||||
productID, code, validGTINType(in.GTINType, code), validPackLevel(in.PackLevel),
|
||||
in.Region, makePrimary, srcArg).
|
||||
Scan(&b.ID, &b.GTIN, &b.GTINType, &b.PackLevel, &b.Region, &b.IsPrimary)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if makePrimary {
|
||||
if _, err := tx.Exec(ctx, "UPDATE product SET gtin=$2 WHERE id=$1", productID, code); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
if _, err := s.recomputeQualityTx(ctx, tx, productID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
_ = s.writeAudit(ctx, actor, "add_barcode", "product", &productID, []string{"gtin"}, nil, b)
|
||||
return &b, nil
|
||||
}
|
||||
|
||||
// DeleteBarcode removes a barcode; if it was primary, another is promoted.
|
||||
func (s *Store) DeleteBarcode(ctx context.Context, productID, barcodeID, actor string) error {
|
||||
tx, err := s.pool.Begin(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer tx.Rollback(ctx)
|
||||
|
||||
var code string
|
||||
var wasPrimary bool
|
||||
err = tx.QueryRow(ctx,
|
||||
"DELETE FROM product_barcode WHERE id=$1 AND product_id=$2 RETURNING gtin, is_primary",
|
||||
barcodeID, productID).Scan(&code, &wasPrimary)
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return ErrNotFound
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if wasPrimary {
|
||||
var newID, newGTIN string
|
||||
e := tx.QueryRow(ctx,
|
||||
"SELECT id, gtin FROM product_barcode WHERE product_id=$1 ORDER BY gtin LIMIT 1", productID).
|
||||
Scan(&newID, &newGTIN)
|
||||
if e == nil {
|
||||
if _, err := tx.Exec(ctx, "UPDATE product_barcode SET is_primary=true WHERE id=$1", newID); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := tx.Exec(ctx, "UPDATE product SET gtin=$2 WHERE id=$1", productID, newGTIN); err != nil {
|
||||
return err
|
||||
}
|
||||
} else if errors.Is(e, pgx.ErrNoRows) {
|
||||
if _, err := tx.Exec(ctx, "UPDATE product SET gtin=NULL WHERE id=$1", productID); err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
return e
|
||||
}
|
||||
}
|
||||
|
||||
if _, err := s.recomputeQualityTx(ctx, tx, productID); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
_ = s.writeAudit(ctx, actor, "delete_barcode", "product", &productID, []string{"gtin"},
|
||||
map[string]string{"gtin": code}, nil)
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetPrimaryBarcode marks one barcode primary and mirrors it to product.gtin.
|
||||
func (s *Store) SetPrimaryBarcode(ctx context.Context, productID, barcodeID, actor string) (*Barcode, error) {
|
||||
tx, err := s.pool.Begin(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer tx.Rollback(ctx)
|
||||
|
||||
var code string
|
||||
err = tx.QueryRow(ctx, "SELECT gtin FROM product_barcode WHERE id=$1 AND product_id=$2", barcodeID, productID).Scan(&code)
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return nil, ErrNotFound
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err := tx.Exec(ctx, "UPDATE product_barcode SET is_primary=false WHERE product_id=$1 AND is_primary", productID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err := tx.Exec(ctx, "UPDATE product_barcode SET is_primary=true WHERE id=$1", barcodeID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err := tx.Exec(ctx, "UPDATE product SET gtin=$2 WHERE id=$1", productID, code); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
_ = s.writeAudit(ctx, actor, "set_primary_barcode", "product", &productID, []string{"gtin"}, nil,
|
||||
map[string]string{"gtin": code})
|
||||
|
||||
bcs, err := s.listBarcodes(ctx, productID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for i := range bcs {
|
||||
if bcs[i].ID == barcodeID {
|
||||
return &bcs[i], nil
|
||||
}
|
||||
}
|
||||
return nil, ErrNotFound
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
// Package gtin validates and normalizes GS1 trade item numbers (GTIN-8/12/13/14).
|
||||
// Only globally-unique GS1 codes are accepted: store-internal / variable-weight /
|
||||
// coupon codes (which are not globally unique) are rejected on purpose.
|
||||
package gtin
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// Validation errors.
|
||||
var (
|
||||
ErrEmpty = errors.New("条码不能为空")
|
||||
ErrFormat = errors.New("条码必须为 8/12/13/14 位数字")
|
||||
ErrCheck = errors.New("条码校验位不正确")
|
||||
ErrRestricted = errors.New("店内码/变量重量码/优惠券码等非全球唯一码,不予收录")
|
||||
)
|
||||
|
||||
// Normalize trims and validates a GTIN, returning the cleaned digit string.
|
||||
// It enforces length, the GS1 mod-10 check digit, and rejects restricted
|
||||
// (non-globally-unique) number ranges.
|
||||
func Normalize(raw string) (string, error) {
|
||||
s := strings.TrimSpace(raw)
|
||||
if s == "" {
|
||||
return "", ErrEmpty
|
||||
}
|
||||
for _, c := range s {
|
||||
if c < '0' || c > '9' {
|
||||
return "", ErrFormat
|
||||
}
|
||||
}
|
||||
switch len(s) {
|
||||
case 8, 12, 13, 14:
|
||||
default:
|
||||
return "", ErrFormat
|
||||
}
|
||||
if !validCheckDigit(s) {
|
||||
return "", ErrCheck
|
||||
}
|
||||
if restricted(s) {
|
||||
return "", ErrRestricted
|
||||
}
|
||||
return s, nil
|
||||
}
|
||||
|
||||
// InferType returns the conventional GTIN type label for a normalized code.
|
||||
func InferType(s string) string {
|
||||
switch len(s) {
|
||||
case 8:
|
||||
return "EAN8"
|
||||
case 12:
|
||||
return "UPC"
|
||||
case 14:
|
||||
return "GTIN14"
|
||||
default:
|
||||
return "EAN13"
|
||||
}
|
||||
}
|
||||
|
||||
// validCheckDigit verifies the trailing GS1 mod-10 check digit. The digit
|
||||
// immediately left of the check digit carries weight 3, then weights alternate.
|
||||
func validCheckDigit(s string) bool {
|
||||
n := len(s)
|
||||
sum := 0
|
||||
for i := 0; i < n-1; i++ {
|
||||
d := int(s[i] - '0')
|
||||
if (n-1-i)%2 == 1 {
|
||||
sum += d * 3
|
||||
} else {
|
||||
sum += d
|
||||
}
|
||||
}
|
||||
check := (10 - (sum % 10)) % 10
|
||||
return check == int(s[n-1]-'0')
|
||||
}
|
||||
|
||||
// restricted reports whether a (length/check-digit valid) code falls in a
|
||||
// number range reserved for non-globally-unique use.
|
||||
func restricted(s string) bool {
|
||||
switch len(s) {
|
||||
case 13:
|
||||
p2 := s[:2]
|
||||
switch {
|
||||
case s[0] == '2': // 20-29 restricted distribution / in-store
|
||||
return true
|
||||
case p2 == "02": // 020-029 variable-measure within a store
|
||||
return true
|
||||
case p2 == "04": // 040-049 restricted circulation within a company
|
||||
return true
|
||||
case p2 == "05": // 050-059 coupons
|
||||
return true
|
||||
case p2 == "98" || p2 == "99": // 980-989/99 coupons & refund receipts
|
||||
return true
|
||||
}
|
||||
case 12: // UPC-A: leading number-system digit
|
||||
switch s[0] {
|
||||
case '2': // in-store / random weight
|
||||
return true
|
||||
case '4': // unrestricted in-store use
|
||||
return true
|
||||
case '5': // coupons
|
||||
return true
|
||||
}
|
||||
case 8: // EAN-8: 0/2 prefixes reserved for in-store use
|
||||
if s[0] == '0' || s[0] == '2' {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
package gtin
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestNormalizeValid(t *testing.T) {
|
||||
cases := []struct{ in, want, typ string }{
|
||||
{" 5449000000996 ", "5449000000996", "EAN13"}, // Coca-Cola EAN-13
|
||||
{"3017624010701", "3017624010701", "EAN13"}, // Nutella EAN-13
|
||||
{"036000291452", "036000291452", "UPC"}, // UPC-A
|
||||
{"96385074", "96385074", "EAN8"}, // EAN-8
|
||||
{"00012345600012", "00012345600012", "GTIN14"},
|
||||
{"6901234567892", "6901234567892", "EAN13"}, // China 690 prefix
|
||||
}
|
||||
for _, c := range cases {
|
||||
got, err := Normalize(c.in)
|
||||
if err != nil {
|
||||
t.Errorf("Normalize(%q) unexpected error: %v", c.in, err)
|
||||
continue
|
||||
}
|
||||
if got != c.want {
|
||||
t.Errorf("Normalize(%q) = %q, want %q", c.in, got, c.want)
|
||||
}
|
||||
if InferType(got) != c.typ {
|
||||
t.Errorf("InferType(%q) = %q, want %q", got, InferType(got), c.typ)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeRejects(t *testing.T) {
|
||||
cases := []struct {
|
||||
in string
|
||||
want error
|
||||
}{
|
||||
{"", ErrEmpty},
|
||||
{"12ab5678", ErrFormat},
|
||||
{"12345", ErrFormat},
|
||||
{"5449000000997", ErrCheck}, // bad check digit
|
||||
{"2012345678903", ErrRestricted}, // 20-29 in-store EAN-13
|
||||
{"0212345678909", ErrRestricted}, // 02x variable measure
|
||||
{"212345678909", ErrRestricted}, // UPC number system 2
|
||||
{"02345673", ErrRestricted}, // EAN-8 in-store
|
||||
}
|
||||
for _, c := range cases {
|
||||
_, err := Normalize(c.in)
|
||||
if err != c.want {
|
||||
t.Errorf("Normalize(%q) error = %v, want %v", c.in, err, c.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -29,6 +29,15 @@ func (s *Store) Ping(ctx context.Context) error {
|
||||
return s.pool.Ping(ctx)
|
||||
}
|
||||
|
||||
// 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"`
|
||||
}
|
||||
|
||||
// Product is the full public view of a product.
|
||||
type Product struct {
|
||||
ID string `json:"id"`
|
||||
@@ -41,6 +50,7 @@ type Product struct {
|
||||
NetContentUnit *string `json:"net_content_unit"`
|
||||
CountryOfOrigin *string `json:"country_of_origin"`
|
||||
QualityScore float64 `json:"quality_score"`
|
||||
Barcodes []Barcode `json:"barcodes"`
|
||||
Nutriments map[string]any `json:"nutriments,omitempty"`
|
||||
NutritionBasis *string `json:"nutrition_basis,omitempty"`
|
||||
NutriScore *string `json:"nutri_score,omitempty"`
|
||||
@@ -49,6 +59,27 @@ type Product struct {
|
||||
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"`
|
||||
@@ -86,16 +117,35 @@ func scanProduct(row pgx.Row) (*Product, error) {
|
||||
return &p, nil
|
||||
}
|
||||
|
||||
// ProductByGTIN looks up an active product by its barcode.
|
||||
// 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.gtin = $1 AND p.status = 'active'", gtin)
|
||||
return scanProduct(row)
|
||||
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, err := scanProduct(row)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if p.Barcodes, err = s.ProductBarcodes(ctx, p.ID); 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)
|
||||
return scanProduct(row)
|
||||
p, err := scanProduct(row)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if p.Barcodes, err = s.ProductBarcodes(ctx, p.ID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return p, nil
|
||||
}
|
||||
|
||||
// SearchProducts performs a fuzzy name search with optional category subtree filter.
|
||||
@@ -104,7 +154,9 @@ func (s *Store) SearchProducts(ctx context.Context, q, category string, limit, o
|
||||
where := "WHERE p.status = 'active'"
|
||||
if q != "" {
|
||||
args = append(args, q)
|
||||
where += " AND p.name ILIKE '%' || $1 || '%'"
|
||||
where += ` AND (p.name ILIKE '%' || $1 || '%'
|
||||
OR EXISTS (SELECT 1 FROM product_barcode pb
|
||||
WHERE pb.product_id = p.id AND pb.gtin ILIKE '%' || $1 || '%'))`
|
||||
}
|
||||
if category != "" {
|
||||
args = append(args, category)
|
||||
|
||||
@@ -27,9 +27,6 @@ _DEFAULT_MIN_INTERVAL = 4.0
|
||||
_API_URL = "https://world.openfoodfacts.org/api/v2/product/{barcode}.json"
|
||||
_SEARCH_URL = "https://world.openfoodfacts.org/api/v2/search"
|
||||
|
||||
# HTTP statuses worth retrying: rate limiting and transient server errors.
|
||||
_RETRY_STATUS = frozenset({429, 500, 502, 503, 504})
|
||||
|
||||
# Fields requested from the search API so a returned product can be transformed
|
||||
# without an extra per-barcode round trip.
|
||||
_SEARCH_FIELDS = (
|
||||
@@ -49,13 +46,9 @@ class OpenFoodFactsAdapter:
|
||||
self,
|
||||
client: httpx.Client | None = None,
|
||||
min_interval: float = _DEFAULT_MIN_INTERVAL,
|
||||
max_retries: int = 4,
|
||||
backoff_base: float = 2.0,
|
||||
) -> None:
|
||||
self._client = client or httpx.Client(headers={"User-Agent": USER_AGENT}, timeout=30.0)
|
||||
self._min_interval = min_interval
|
||||
self._max_retries = max_retries
|
||||
self._backoff_base = backoff_base
|
||||
self._last_call = 0.0
|
||||
|
||||
def _throttle(self) -> None:
|
||||
@@ -65,49 +58,11 @@ class OpenFoodFactsAdapter:
|
||||
time.sleep(wait)
|
||||
self._last_call = time.monotonic()
|
||||
|
||||
def _get(self, url: str, params: dict | None = None) -> httpx.Response:
|
||||
"""GET with throttling and retry/backoff on transient errors.
|
||||
|
||||
Retries on connection/timeout errors and on retryable HTTP statuses
|
||||
(429 and 5xx, which OFF returns intermittently when overloaded), using
|
||||
exponential backoff that honours a ``Retry-After`` header when present.
|
||||
"""
|
||||
last_exc: Exception | None = None
|
||||
for attempt in range(self._max_retries + 1):
|
||||
self._throttle()
|
||||
try:
|
||||
resp = self._client.get(url, params=params)
|
||||
except httpx.TransportError as exc:
|
||||
last_exc = exc
|
||||
else:
|
||||
if resp.status_code < 400 or resp.status_code not in _RETRY_STATUS:
|
||||
resp.raise_for_status()
|
||||
return resp
|
||||
last_exc = httpx.HTTPStatusError(
|
||||
f"retryable status {resp.status_code}", request=resp.request, response=resp
|
||||
)
|
||||
if attempt < self._max_retries:
|
||||
retry_after = self._retry_after(last_exc)
|
||||
time.sleep(retry_after if retry_after is not None else self._backoff_base**attempt)
|
||||
assert last_exc is not None
|
||||
raise last_exc
|
||||
|
||||
@staticmethod
|
||||
def _retry_after(exc: Exception | None) -> float | None:
|
||||
resp = getattr(exc, "response", None)
|
||||
if resp is None:
|
||||
return None
|
||||
value = resp.headers.get("Retry-After")
|
||||
if not value:
|
||||
return None
|
||||
try:
|
||||
return float(value)
|
||||
except ValueError:
|
||||
return None
|
||||
|
||||
def fetch_barcode(self, barcode: str) -> dict | None:
|
||||
"""Fetch a single product by barcode; return the raw `product` dict."""
|
||||
resp = self._get(_API_URL.format(barcode=barcode))
|
||||
self._throttle()
|
||||
resp = self._client.get(_API_URL.format(barcode=barcode))
|
||||
resp.raise_for_status()
|
||||
payload = resp.json()
|
||||
if payload.get("status") != 1:
|
||||
return None
|
||||
@@ -136,7 +91,8 @@ class OpenFoodFactsAdapter:
|
||||
``last_modified_t`` they processed as the next watermark.
|
||||
"""
|
||||
for page in range(1, max_pages + 1):
|
||||
resp = self._get(
|
||||
self._throttle()
|
||||
resp = self._client.get(
|
||||
_SEARCH_URL,
|
||||
params={
|
||||
"fields": _SEARCH_FIELDS,
|
||||
@@ -145,6 +101,7 @@ class OpenFoodFactsAdapter:
|
||||
"page_size": page_size,
|
||||
},
|
||||
)
|
||||
resp.raise_for_status()
|
||||
products = resp.json().get("products") or []
|
||||
if not products:
|
||||
return
|
||||
@@ -157,60 +114,6 @@ class OpenFoodFactsAdapter:
|
||||
if reached_old or len(products) < page_size:
|
||||
return
|
||||
|
||||
def fetch_by_country(
|
||||
self,
|
||||
country: str,
|
||||
*,
|
||||
page_size: int = 100,
|
||||
max_pages: int = 10,
|
||||
sort_by: str = "unique_scans_n",
|
||||
) -> Iterator[dict]:
|
||||
"""Yield products sold in ``country`` (an OFF ``countries_tags_en`` slug).
|
||||
|
||||
Used to seed a market-specific catalogue (e.g. ``china``). Results are
|
||||
sorted by ``sort_by`` (default ``unique_scans_n`` so the most-scanned,
|
||||
best-known products come first) and de-duplicated across pages, since
|
||||
OFF's popularity ordering is not stable between page requests.
|
||||
"""
|
||||
seen: set[str] = set()
|
||||
for page in range(1, max_pages + 1):
|
||||
resp = self._get(
|
||||
_SEARCH_URL,
|
||||
params={
|
||||
"fields": _SEARCH_FIELDS,
|
||||
"countries_tags_en": country,
|
||||
"sort_by": sort_by,
|
||||
"page": page,
|
||||
"page_size": page_size,
|
||||
},
|
||||
)
|
||||
products = resp.json().get("products") or []
|
||||
if not products:
|
||||
return
|
||||
new_on_page = 0
|
||||
for prod in products:
|
||||
code = str(prod.get("code") or "")
|
||||
if code and code in seen:
|
||||
continue
|
||||
if code:
|
||||
seen.add(code)
|
||||
new_on_page += 1
|
||||
yield prod
|
||||
if len(products) < page_size or new_on_page == 0:
|
||||
return
|
||||
|
||||
|
||||
def is_cn_gs1(code: str | None) -> bool:
|
||||
"""Return True for a GS1 China company prefix (barcodes starting 690-699).
|
||||
|
||||
These identify products registered with GS1 China, i.e. genuinely domestic
|
||||
items, as opposed to imported goods merely tagged as sold in China.
|
||||
"""
|
||||
if not code:
|
||||
return False
|
||||
code = code.strip()
|
||||
return len(code) >= 3 and code[:2] == "69" and code[2].isdigit()
|
||||
|
||||
|
||||
def read_dump(path: str | Path) -> Iterator[dict]:
|
||||
"""Yield raw product records from an OFF JSONL dump file.
|
||||
|
||||
@@ -7,7 +7,6 @@ source with field-level provenance in `product_source`.
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
from typing import Any
|
||||
|
||||
@@ -19,8 +18,6 @@ from opengoods.etl.quality import update_quality
|
||||
|
||||
OFF_HOMEPAGE = "https://world.openfoodfacts.org"
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def default_dsn() -> str:
|
||||
return os.environ.get(
|
||||
@@ -199,25 +196,6 @@ def load_record(conn: psycopg.Connection, rec: dict[str, Any], source_id: str, r
|
||||
return product_id
|
||||
|
||||
|
||||
def load_record_safe(
|
||||
conn: psycopg.Connection, rec: dict[str, Any], source_id: str, raw: dict
|
||||
) -> bool:
|
||||
"""Load one record inside a savepoint.
|
||||
|
||||
On success the record's writes stay in the surrounding transaction. On any
|
||||
error, only this record's writes are rolled back (to the savepoint) and the
|
||||
batch continues, so a single malformed source record cannot abort a large
|
||||
import. Returns True if loaded, False if skipped due to an error.
|
||||
"""
|
||||
try:
|
||||
with conn.transaction():
|
||||
load_record(conn, rec, source_id, raw)
|
||||
return True
|
||||
except Exception as exc: # noqa: BLE001 - per-record isolation is intentional
|
||||
logger.warning("skipping record gtin=%s: %s", rec.get("gtin"), exc)
|
||||
return False
|
||||
|
||||
|
||||
def _jsonable(raw: dict) -> dict:
|
||||
"""Drop values that are not JSON-serializable from a raw record."""
|
||||
try:
|
||||
|
||||
@@ -78,14 +78,6 @@ def map_category(raw: dict) -> str | None:
|
||||
return None
|
||||
|
||||
|
||||
def _clamp(value: str | None, max_len: int) -> str | None:
|
||||
"""Trim a string to fit a bounded DB column; external data length varies."""
|
||||
if value is None:
|
||||
return None
|
||||
value = value.strip()
|
||||
return value[:max_len] or None
|
||||
|
||||
|
||||
def _clean_tags(tags: list[str] | None, prefix: str = "") -> list[str]:
|
||||
out: list[str] = []
|
||||
for t in tags or []:
|
||||
@@ -151,16 +143,16 @@ def transform(raw: dict) -> dict | None:
|
||||
"brand": brand,
|
||||
"category_path": map_category(raw),
|
||||
"net_content_value": net_value,
|
||||
"net_content_unit": _clamp(net_unit, 16),
|
||||
"net_content_unit": net_unit,
|
||||
"net_content_canonical": net_canonical,
|
||||
"country_of_origin": _clamp((raw.get("countries") or "").split(",")[0].strip() or None, 64),
|
||||
"country_of_origin": (raw.get("countries") or "").split(",")[0].strip() or None,
|
||||
"food": {
|
||||
"ingredients_text": raw.get("ingredients_text") or None,
|
||||
"allergens": _clean_tags(raw.get("allergens_tags")),
|
||||
"additives": _clean_tags(raw.get("additives_tags")),
|
||||
"nutriments": transform_nutriments(raw.get("nutriments") or {}),
|
||||
"nutrition_basis": "per_100g",
|
||||
"serving_size": _clamp(raw.get("serving_size") or None, 32),
|
||||
"serving_size": raw.get("serving_size") or None,
|
||||
"nutri_score": (raw.get("nutriscore_grade") or "").upper()[:1] or None,
|
||||
},
|
||||
"image_url": raw.get("image_front_url") or raw.get("image_url") or None,
|
||||
|
||||
@@ -7,11 +7,6 @@ Usage:
|
||||
# from a downloaded OFF JSONL dump (optionally .gz), limited to N records
|
||||
python -m opengoods.jobs.seed_off --dump products.jsonl.gz --limit 1000
|
||||
|
||||
# market-focused: the most-scanned products sold in China, restricted to
|
||||
# genuine GS1-China (69x) barcodes
|
||||
python -m opengoods.jobs.seed_off --country china --domestic-only \
|
||||
--max-pages 20 --limit 1000
|
||||
|
||||
The OFF read API is rate-limited client-side; for large imports use a dump.
|
||||
"""
|
||||
|
||||
@@ -23,38 +18,25 @@ from collections.abc import Iterator
|
||||
|
||||
import psycopg
|
||||
|
||||
from opengoods.adapters.openfoodfacts import (
|
||||
OpenFoodFactsAdapter,
|
||||
is_cn_gs1,
|
||||
read_dump,
|
||||
)
|
||||
from opengoods.etl.load import default_dsn, ensure_source, load_record_safe
|
||||
from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter, read_dump
|
||||
from opengoods.etl.load import default_dsn, ensure_source, load_record
|
||||
from opengoods.etl.transform import transform
|
||||
|
||||
|
||||
def _raw_records(args: argparse.Namespace) -> Iterator[dict]:
|
||||
if args.dump:
|
||||
records: Iterator[dict] = read_dump(args.dump)
|
||||
elif args.country:
|
||||
adapter = OpenFoodFactsAdapter(min_interval=args.min_interval)
|
||||
records = adapter.fetch_by_country(
|
||||
args.country, page_size=args.page_size, max_pages=args.max_pages
|
||||
)
|
||||
records = read_dump(args.dump)
|
||||
else:
|
||||
adapter = OpenFoodFactsAdapter(min_interval=args.min_interval)
|
||||
records = adapter.fetch(args.barcodes)
|
||||
yielded = 0
|
||||
for rec in records:
|
||||
if args.domestic_only and not is_cn_gs1(str(rec.get("code") or "")):
|
||||
continue
|
||||
if args.limit and yielded >= args.limit:
|
||||
for i, rec in enumerate(records):
|
||||
if args.limit and i >= args.limit:
|
||||
break
|
||||
yielded += 1
|
||||
yield rec
|
||||
|
||||
|
||||
def run(args: argparse.Namespace) -> int:
|
||||
loaded = skipped = errored = 0
|
||||
loaded = skipped = 0
|
||||
with psycopg.connect(args.dsn, autocommit=False) as conn:
|
||||
source_id = ensure_source(conn)
|
||||
for raw in _raw_records(args):
|
||||
@@ -62,12 +44,10 @@ def run(args: argparse.Namespace) -> int:
|
||||
if rec is None:
|
||||
skipped += 1
|
||||
continue
|
||||
if load_record_safe(conn, rec, source_id, raw):
|
||||
loaded += 1
|
||||
else:
|
||||
errored += 1
|
||||
load_record(conn, rec, source_id, raw)
|
||||
loaded += 1
|
||||
conn.commit()
|
||||
print(f"loaded={loaded} skipped={skipped} errored={errored}")
|
||||
print(f"loaded={loaded} skipped={skipped}")
|
||||
return 0
|
||||
|
||||
|
||||
@@ -76,18 +56,7 @@ def main(argv: list[str] | None = None) -> int:
|
||||
src = parser.add_mutually_exclusive_group(required=True)
|
||||
src.add_argument("--barcodes", nargs="+", help="barcodes to fetch via the OFF API")
|
||||
src.add_argument("--dump", help="path to an OFF JSONL dump (.jsonl or .jsonl.gz)")
|
||||
src.add_argument(
|
||||
"--country",
|
||||
help="OFF countries_tags_en slug to seed from, e.g. 'china' (most-scanned first)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--domestic-only",
|
||||
action="store_true",
|
||||
help="keep only genuine GS1-China (69x) barcodes; drop imported goods",
|
||||
)
|
||||
parser.add_argument("--limit", type=int, default=0, help="max records to load (0 = all)")
|
||||
parser.add_argument("--page-size", type=int, default=100, help="search page size")
|
||||
parser.add_argument("--max-pages", type=int, default=10, help="max search pages (country mode)")
|
||||
parser.add_argument("--min-interval", type=float, default=4.0, help="API throttle seconds")
|
||||
parser.add_argument("--dsn", default=default_dsn(), help="PostgreSQL DSN")
|
||||
return run(parser.parse_args(argv))
|
||||
|
||||
@@ -17,14 +17,14 @@ import sys
|
||||
import psycopg
|
||||
|
||||
from opengoods.adapters.openfoodfacts import SOURCE_NAME, OpenFoodFactsAdapter
|
||||
from opengoods.etl.load import default_dsn, ensure_source, load_record_safe
|
||||
from opengoods.etl.load import default_dsn, ensure_source, load_record
|
||||
from opengoods.etl.state import get_watermark, set_watermark
|
||||
from opengoods.etl.transform import transform
|
||||
|
||||
|
||||
def run(args: argparse.Namespace) -> int:
|
||||
adapter = OpenFoodFactsAdapter(min_interval=args.min_interval)
|
||||
loaded = skipped = errored = 0
|
||||
loaded = skipped = 0
|
||||
high_watermark = 0
|
||||
with psycopg.connect(args.dsn, autocommit=False) as conn:
|
||||
source_id = ensure_source(conn)
|
||||
@@ -38,21 +38,16 @@ def run(args: argparse.Namespace) -> int:
|
||||
if rec is None:
|
||||
skipped += 1
|
||||
continue
|
||||
if load_record_safe(conn, rec, source_id, raw):
|
||||
loaded += 1
|
||||
else:
|
||||
errored += 1
|
||||
load_record(conn, rec, source_id, raw)
|
||||
loaded += 1
|
||||
set_watermark(
|
||||
conn,
|
||||
SOURCE_NAME,
|
||||
high_watermark,
|
||||
stats={"loaded": loaded, "skipped": skipped, "errored": errored, "since": since},
|
||||
stats={"loaded": loaded, "skipped": skipped, "since": since},
|
||||
)
|
||||
conn.commit()
|
||||
print(
|
||||
f"since={since} loaded={loaded} skipped={skipped} "
|
||||
f"errored={errored} watermark={high_watermark}"
|
||||
)
|
||||
print(f"since={since} loaded={loaded} skipped={skipped} watermark={high_watermark}")
|
||||
return 0
|
||||
|
||||
|
||||
|
||||
@@ -11,7 +11,7 @@ from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from opengoods.etl.load import default_dsn, ensure_source, load_record, load_record_safe
|
||||
from opengoods.etl.load import default_dsn, ensure_source, load_record
|
||||
from opengoods.etl.transform import transform
|
||||
|
||||
psycopg = pytest.importorskip("psycopg")
|
||||
@@ -59,25 +59,3 @@ def test_load_record_roundtrip(conn):
|
||||
assert prov[0] >= 1
|
||||
|
||||
conn.rollback() # keep the test DB clean
|
||||
|
||||
|
||||
def test_load_record_safe_isolates_bad_record(conn):
|
||||
source_id = ensure_source(conn)
|
||||
|
||||
# Unique gtin so the good record is a fresh INSERT, not an upsert/update.
|
||||
good = transform(FIXTURE)
|
||||
good["gtin"] = "4006381333931"
|
||||
assert load_record_safe(conn, good, source_id, FIXTURE) is True
|
||||
after_good = conn.execute("SELECT count(*) FROM product").fetchone()[0]
|
||||
|
||||
# A record whose name violates NOT NULL must not abort the batch.
|
||||
bad = dict(good)
|
||||
bad["gtin"] = "5000112637922"
|
||||
bad["name"] = None
|
||||
assert load_record_safe(conn, bad, source_id, {}) is False
|
||||
|
||||
# The good record survived the bad one's rollback-to-savepoint.
|
||||
after_bad = conn.execute("SELECT count(*) FROM product").fetchone()[0]
|
||||
assert after_bad == after_good
|
||||
|
||||
conn.rollback() # keep the test DB clean
|
||||
|
||||
@@ -1,53 +0,0 @@
|
||||
import httpx
|
||||
|
||||
from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter, is_cn_gs1
|
||||
|
||||
|
||||
def _adapter(handler, **kwargs):
|
||||
client = httpx.Client(transport=httpx.MockTransport(handler))
|
||||
return OpenFoodFactsAdapter(client=client, min_interval=0, backoff_base=0, **kwargs)
|
||||
|
||||
|
||||
def test_is_cn_gs1():
|
||||
assert is_cn_gs1("6901234567892") # GS1 China prefix
|
||||
assert is_cn_gs1("690")
|
||||
assert not is_cn_gs1("3017624010701") # France
|
||||
assert not is_cn_gs1("5449000000996") # Belgium
|
||||
assert not is_cn_gs1("")
|
||||
assert not is_cn_gs1(None)
|
||||
assert not is_cn_gs1("69") # too short to carry a prefix digit
|
||||
|
||||
|
||||
def test_fetch_by_country_passes_filter_and_paginates():
|
||||
seen_params = []
|
||||
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
seen_params.append(dict(request.url.params))
|
||||
page = int(request.url.params.get("page", "1"))
|
||||
if page == 1:
|
||||
return httpx.Response(
|
||||
200,
|
||||
json={"products": [{"code": "6901"}, {"code": "6902"}]},
|
||||
)
|
||||
return httpx.Response(200, json={"products": []})
|
||||
|
||||
adapter = _adapter(handler)
|
||||
out = list(adapter.fetch_by_country("china", page_size=2, max_pages=5))
|
||||
assert [p["code"] for p in out] == ["6901", "6902"]
|
||||
assert seen_params[0]["countries_tags_en"] == "china"
|
||||
assert seen_params[0]["sort_by"] == "unique_scans_n"
|
||||
|
||||
|
||||
def test_fetch_by_country_dedupes_across_pages():
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
page = int(request.url.params.get("page", "1"))
|
||||
if page == 1:
|
||||
return httpx.Response(200, json={"products": [{"code": "6901"}, {"code": "6902"}]})
|
||||
if page == 2:
|
||||
# OFF popularity ordering is unstable; a repeat appears on page 2
|
||||
return httpx.Response(200, json={"products": [{"code": "6902"}, {"code": "6903"}]})
|
||||
return httpx.Response(200, json={"products": []})
|
||||
|
||||
adapter = _adapter(handler)
|
||||
out = [p["code"] for p in adapter.fetch_by_country("china", page_size=2, max_pages=5)]
|
||||
assert out == ["6901", "6902", "6903"]
|
||||
@@ -1,59 +0,0 @@
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter
|
||||
|
||||
|
||||
def _adapter(handler, **kwargs):
|
||||
client = httpx.Client(transport=httpx.MockTransport(handler))
|
||||
return OpenFoodFactsAdapter(client=client, min_interval=0, backoff_base=0, **kwargs)
|
||||
|
||||
|
||||
def test_retries_transient_5xx_then_succeeds():
|
||||
calls = {"n": 0}
|
||||
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
calls["n"] += 1
|
||||
if calls["n"] < 3:
|
||||
return httpx.Response(503)
|
||||
return httpx.Response(200, json={"status": 1, "product": {"code": "x"}})
|
||||
|
||||
adapter = _adapter(handler, max_retries=4)
|
||||
assert adapter.fetch_barcode("x") == {"code": "x"}
|
||||
assert calls["n"] == 3 # two 503s retried, third succeeds
|
||||
|
||||
|
||||
def test_retries_exhausted_raises():
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
return httpx.Response(503)
|
||||
|
||||
adapter = _adapter(handler, max_retries=2)
|
||||
with pytest.raises(httpx.HTTPStatusError):
|
||||
adapter.fetch_barcode("x")
|
||||
|
||||
|
||||
def test_non_retryable_4xx_not_retried():
|
||||
calls = {"n": 0}
|
||||
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
calls["n"] += 1
|
||||
return httpx.Response(404)
|
||||
|
||||
adapter = _adapter(handler, max_retries=4)
|
||||
with pytest.raises(httpx.HTTPStatusError):
|
||||
adapter.fetch_barcode("x")
|
||||
assert calls["n"] == 1 # 404 is not retried
|
||||
|
||||
|
||||
def test_retries_connection_error_then_succeeds():
|
||||
calls = {"n": 0}
|
||||
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
calls["n"] += 1
|
||||
if calls["n"] == 1:
|
||||
raise httpx.ConnectError("boom")
|
||||
return httpx.Response(200, json={"status": 1, "product": {"code": "y"}})
|
||||
|
||||
adapter = _adapter(handler, max_retries=4)
|
||||
assert adapter.fetch_barcode("y") == {"code": "y"}
|
||||
assert calls["n"] == 2
|
||||
@@ -65,14 +65,3 @@ def test_transform_full_record():
|
||||
|
||||
def test_transform_drops_unnamed():
|
||||
assert transform({"code": "0000000000000"}) is None
|
||||
|
||||
|
||||
def test_transform_clamps_long_serving_size():
|
||||
raw = {
|
||||
"code": "3017624010701",
|
||||
"product_name": "X",
|
||||
"serving_size": "1 portion (30 g) / 1 portion (30 g) / extra long descriptive text here",
|
||||
}
|
||||
rec = transform(raw)
|
||||
assert rec is not None
|
||||
assert len(rec["food"]["serving_size"]) == 32
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
DROP TABLE IF EXISTS product_barcode;
|
||||
@@ -0,0 +1,36 @@
|
||||
-- Multi-barcode support: one product can carry many GS1 barcodes
|
||||
-- (consumer unit EAN-13/UPC, case ITF-14, regional re-labels, etc.).
|
||||
-- product.gtin is kept as the denormalized "primary" barcode for
|
||||
-- backward compatibility and is mirrored from the is_primary row here.
|
||||
|
||||
CREATE TABLE product_barcode (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
product_id UUID NOT NULL REFERENCES product(id) ON DELETE CASCADE,
|
||||
gtin VARCHAR(14) NOT NULL,
|
||||
gtin_type VARCHAR(8) NOT NULL DEFAULT 'EAN13',
|
||||
pack_level VARCHAR(8) NOT NULL DEFAULT 'each',
|
||||
region VARCHAR(8),
|
||||
is_primary BOOLEAN NOT NULL DEFAULT false,
|
||||
source_id UUID REFERENCES source(id),
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
CONSTRAINT product_barcode_type_chk CHECK (gtin_type IN ('EAN8','UPC','EAN13','ITF14','GTIN14')),
|
||||
CONSTRAINT product_barcode_pack_chk CHECK (pack_level IN ('each','case','pallet'))
|
||||
);
|
||||
|
||||
-- A barcode is globally unique: one code maps to exactly one product.
|
||||
CREATE UNIQUE INDEX idx_product_barcode_gtin ON product_barcode (gtin);
|
||||
CREATE INDEX idx_product_barcode_product ON product_barcode (product_id);
|
||||
-- At most one primary barcode per product.
|
||||
CREATE UNIQUE INDEX idx_product_barcode_primary ON product_barcode (product_id) WHERE is_primary;
|
||||
|
||||
-- Backfill: lift each product's existing gtin into the new table as primary.
|
||||
INSERT INTO product_barcode (product_id, gtin, gtin_type, pack_level, is_primary)
|
||||
SELECT id, gtin,
|
||||
CASE WHEN length(gtin) = 8 THEN 'EAN8'
|
||||
WHEN length(gtin) = 12 THEN 'UPC'
|
||||
WHEN length(gtin) = 14 THEN 'GTIN14'
|
||||
ELSE 'EAN13' END,
|
||||
'each', true
|
||||
FROM product
|
||||
WHERE gtin IS NOT NULL AND gtin <> ''
|
||||
ON CONFLICT (gtin) DO NOTHING;
|
||||
@@ -60,6 +60,29 @@ export default function ProductView({ id, onBack }: { id: string; onBack: () =>
|
||||
<div className="mt-4">
|
||||
<Row label="品牌" value={p.brand} />
|
||||
<Row label="条码 (GTIN)" value={p.gtin} />
|
||||
{(() => {
|
||||
const others = (p.barcodes || []).filter(
|
||||
(b) => !b.is_primary && b.gtin !== p.gtin,
|
||||
);
|
||||
return others.length ? (
|
||||
<Row
|
||||
label="其他条码"
|
||||
value={
|
||||
<div className="flex flex-wrap gap-1.5">
|
||||
{others.map((b) => (
|
||||
<span
|
||||
key={b.gtin}
|
||||
className="font-mono text-xs bg-gray-100 text-gray-600 rounded px-1.5 py-0.5"
|
||||
title={`${b.gtin_type} · ${b.pack_level}`}
|
||||
>
|
||||
{b.gtin}
|
||||
</span>
|
||||
))}
|
||||
</div>
|
||||
}
|
||||
/>
|
||||
) : null;
|
||||
})()}
|
||||
<Row label="品类" value={p.category_path} />
|
||||
<Row
|
||||
label="净含量"
|
||||
|
||||
@@ -6,9 +6,18 @@ export interface ProductSummary {
|
||||
category_path: string | null;
|
||||
}
|
||||
|
||||
export interface Barcode {
|
||||
gtin: string;
|
||||
gtin_type: string;
|
||||
pack_level: string;
|
||||
region: string | null;
|
||||
is_primary: boolean;
|
||||
}
|
||||
|
||||
export interface Product {
|
||||
id: string;
|
||||
gtin: string | null;
|
||||
barcodes?: Barcode[] | null;
|
||||
name: string;
|
||||
brand: string | null;
|
||||
category_path: string | null;
|
||||
|
||||
Reference in New Issue
Block a user