Compare commits

..

13 Commits

Author SHA1 Message Date
oyaegeli98668 5bfb2d3a5f chore: 站名改为「天工商品档案公共仓」(公开站 + 后台标题/页头/页脚/API 文案)
CI / Go (api) (pull_request) Failing after 15s
CI / Python (ingestion) (pull_request) Successful in 7s
CI / Migrations (postgres) (pull_request) Failing after 19s
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 05:37:23 +00:00
oyaegeli98668 8909ddcb69 chore(public): 页脚添加 ICP 备案号 鲁ICP备2025185218号-6
CI / Go (api) (pull_request) Failing after 20s
CI / Python (ingestion) (pull_request) Successful in 8s
CI / Migrations (postgres) (pull_request) Failing after 16s
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 05:33:58 +00:00
oyaegeli98668 2d95052bd9 feat(public): 首页新增 API 调用说明页(只读端点/参数/示例/免责声明)
CI / Python (ingestion) (pull_request) Successful in 6s
CI / Go (api) (pull_request) Failing after 17s
CI / Migrations (postgres) (pull_request) Failing after 20s
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 05:26:53 +00:00
oyaegeli98668 c9a4404052 feat: 公开首页(搜索+商品详情) + 好心人投稿 + 后台审核收纳
CI / Go (api) (pull_request) Failing after 20s
CI / Python (ingestion) (pull_request) Successful in 8s
CI / Migrations (postgres) (pull_request) Failing after 17s
- 公开前端 SPA(根路径 /):首页大搜索框、检索结果、只读商品详情、贡献档案表单
- 公开写入端点 POST /api/public/submissions(无需登录,基础频率限流),投稿进入 submission 待审核队列,不直接写 product
- 迁移 0006:submission 投稿表 + community 来源(trust=0.50)
- 后台审核队列:列表(待审核/已通过/已驳回) → 查看投稿 → 通过(创建/补全商品 + 记 source=community + 字段级溯源 + 审计 + 重算质量分) / 驳回(记原因)
- 公开只读 api 服务内嵌公开 SPA;Dockerfile.prod 增加 node 构建阶段 + 内嵌 dist

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 05:11:07 +00:00
oyaegeli98668 bad771e172 build(admin): npm 走 npmmirror 镜像以适配部署环境网络
CI / Go (api) (pull_request) Failing after 19s
CI / Python (ingestion) (pull_request) Successful in 7s
CI / Migrations (postgres) (pull_request) Failing after 21s
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 02:54:32 +00:00
oyaegeli98668 d90a539e6b feat(admin): 运营后台(登录/查看/审核编辑/补全)+ 写入API + 审计留痕
CI / Go (api) (pull_request) Failing after 18s
CI / Python (ingestion) (pull_request) Successful in 7s
CI / Migrations (postgres) (pull_request) Failing after 18s
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 02:53:30 +00:00
lixu c272b2c5d7 Merge pull request 'chore(deploy): 生产 docker compose 部署配置' (#3) from devin/1781922421-prod-deploy into main
CI / Go (api) (push) Failing after 20s
CI / Python (ingestion) (push) Successful in 6s
CI / Migrations (postgres) (push) Failing after 19s
2026-06-20 10:32:22 +08:00
oyaegeli98668 6777461268 chore(deploy): 生产 docker compose + scratch Dockerfile + 部署文档
CI / Go (api) (pull_request) Failing after 41s
CI / Python (ingestion) (pull_request) Successful in 43s
CI / Migrations (postgres) (pull_request) Failing after 17s
- docker-compose.prod.yml:仅 api 绑 127.0.0.1:8120,DB/redis/minio 不暴露,restart=unless-stopped,凭据走 .env
- api/Dockerfile.prod:runtime 改用 scratch + 拷贝 ca-certs(gcr.io/distroless 不可达环境),Go 模块走 goproxy.cn
- .env.example / docs/deploy.md

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 02:27:09 +00:00
lixu e3e5c7979c ci: 修复 Gitea Actions 数据库连接(用服务名 postgres) (#2)
CI / Go (api) (push) Successful in 6s
CI / Python (ingestion) (push) Successful in 6s
CI / Migrations (postgres) (push) Successful in 16s
2026-06-19 17:47:41 +08:00
rosemariejebbjtxbfp 19b7c43f37 ci: connect to postgres via service hostname for Gitea Actions
CI / Go (api) (pull_request) Successful in 30s
CI / Python (ingestion) (pull_request) Successful in 7s
CI / Migrations (postgres) (pull_request) Successful in 16s
Gitea Actions runs jobs inside a container, so a service container is reachable by its service name (postgres), not localhost; localhost:5432 yields connection refused. Also drop the unnecessary 5432:5432 host port mapping that caused 'port already allocated' collisions.

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-19 08:54:18 +00:00
lixu 57537c880d M4: ingestion management (#1)
CI / Go (api) (push) Failing after 14s
CI / Migrations (postgres) (push) Failing after 33s
CI / Python (ingestion) (push) Failing after 13m10s
Incremental OFF + GS1 supplement + dedup/conflict + quality scoring + scheduler
2026-06-08 17:40:26 +08:00
John Doe a35bcd6647 M4: ingestion management (incremental, GS1 supplement, dedup/conflict, quality, scheduler)
CI / Go (api) (pull_request) Has been cancelled
CI / Python (ingestion) (pull_request) Has been cancelled
CI / Migrations (postgres) (pull_request) Has been cancelled
- OFF incremental fetch via search API + persistent watermark (ingest_state, migration 0004)
- GS1 barcode supplement adapter (offline mapping + GS1-style API) filling only gaps with field-level provenance
- Non-GTIN dedup with canonical selection + merge_log; field-level conflict resolution (source trust > recency)
- Quality scoring (0.4 completeness + 0.3 source trust + 0.2 multi-source + 0.1 freshness) wired into load/merge
- Jobs: update_off, dedup, schedule; docs/ingestion-management.md
- 19 new tests (pure + DB-integration), ruff clean

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-08 09:27:42 +00:00
lixu e750501b44 feat(M3): 只读 API 端点实现
CI / Go (api) (push) Has been cancelled
CI / Python (ingestion) (push) Has been cancelled
CI / Migrations (postgres) (push) Has been cancelled
- store: pgx 只读数据访问层(productByGTIN/ByID/search/nutriments/msrp/brands/categories/source)
- handler: 真实查询替换 501 占位, 统一分页 + 错误信封, MSRP 带免责声明无购买入口
- main: pgxpool 连接池接线
- search: 名称模糊 + 分类子树过滤(ltree <@)
- 测试: healthz/pageParams 单测 + DB-backed handler 集成测试(无库自动跳过)
- CI: Go job 增加 postgres service + migrate up, 实跑 DB 测试
- 依赖: pgx v5.7.2 (固定到兼容 go1.23 的版本)
- 本地实跑: 8 个端点对真实 OFF 数据返回正确(barcode/search/nutriments/msrp/brands/categories/source/404)

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-08 07:14:18 +00:00
92 changed files with 12577 additions and 32 deletions
+14
View File
@@ -0,0 +1,14 @@
.git
**/node_modules
admin-frontend/dist
public-frontend/dist
api/server
*.test
*.out
__pycache__
.venv
.pytest_cache
.ruff_cache
.env
.env.*
!.env.example
+12
View File
@@ -0,0 +1,12 @@
# Copy to .env and fill in real values before running docker-compose.prod.yml.
# Used by docker-compose.prod.yml for production deployment.
POSTGRES_USER=opengoods
POSTGRES_PASSWORD=change-me
POSTGRES_DB=opengoods
MINIO_ROOT_USER=opengoods
MINIO_ROOT_PASSWORD=change-me
# Admin console (served at /ping). Set a strong password and a random JWT secret.
GOODS_ADMIN_USER=admin
GOODS_ADMIN_PASSWORD=change-me
GOODS_ADMIN_JWT_SECRET=change-me-to-a-long-random-string
+18 -3
View File
@@ -12,12 +12,29 @@ jobs:
defaults:
run:
working-directory: api
services:
postgres:
image: postgres:16-alpine
env:
POSTGRES_USER: opengoods
POSTGRES_PASSWORD: opengoods
POSTGRES_DB: opengoods
options: >-
--health-cmd "pg_isready -U opengoods"
--health-interval 5s --health-timeout 5s --health-retries 10
env:
OPENGOODS_DATABASE_URL: postgres://opengoods:opengoods@postgres:5432/opengoods?sslmode=disable
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version: "1.23"
cache-dependency-path: api/go.sum
- name: Apply migrations
working-directory: .
run: |
go install -tags 'postgres' github.com/golang-migrate/migrate/v4/cmd/migrate@v4.18.1
migrate -path migrations -database "$OPENGOODS_DATABASE_URL" up
- name: Verify gofmt
run: test -z "$(gofmt -l .)"
- run: go vet ./...
@@ -54,13 +71,11 @@ jobs:
POSTGRES_USER: opengoods
POSTGRES_PASSWORD: opengoods
POSTGRES_DB: opengoods
ports:
- "5432:5432"
options: >-
--health-cmd "pg_isready -U opengoods"
--health-interval 5s --health-timeout 5s --health-retries 10
env:
DBURL: postgres://opengoods:opengoods@localhost:5432/opengoods?sslmode=disable
DBURL: postgres://opengoods:opengoods@postgres:5432/opengoods?sslmode=disable
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
+9
View File
@@ -18,6 +18,15 @@ dist/
.env.*
!.env.example
# Node / admin frontend
node_modules/
# Keep the embedded SPA placeholders (real builds are injected during Docker build)
!api/internal/adminweb/dist/
!api/internal/adminweb/dist/index.html
!api/internal/publicweb/dist/
!api/internal/publicweb/dist/index.html
# OS / editors
.DS_Store
*.swp
+12
View File
@@ -0,0 +1,12 @@
<!doctype html>
<html lang="zh">
<head>
<meta charset="UTF-8" />
<meta name="viewport" content="width=device-width, initial-scale=1.0" />
<title>天工商品档案公共仓 · 管理后台</title>
</head>
<body>
<div id="root"></div>
<script type="module" src="/src/main.tsx"></script>
</body>
</html>
+2690
View File
File diff suppressed because it is too large Load Diff
+26
View File
@@ -0,0 +1,26 @@
{
"name": "opengoods-admin-frontend",
"private": true,
"version": "1.0.0",
"type": "module",
"scripts": {
"dev": "vite",
"build": "tsc && vite build",
"preview": "vite preview"
},
"dependencies": {
"react": "^18.2.0",
"react-dom": "^18.2.0",
"lucide-react": "^0.344.0"
},
"devDependencies": {
"@types/react": "^18.2.55",
"@types/react-dom": "^18.2.19",
"@vitejs/plugin-react": "^4.2.1",
"autoprefixer": "^10.4.17",
"postcss": "^8.4.35",
"tailwindcss": "^3.4.1",
"typescript": "^5.3.3",
"vite": "^5.1.0"
}
}
+6
View File
@@ -0,0 +1,6 @@
export default {
plugins: {
tailwindcss: {},
autoprefixer: {},
},
};
+125
View File
@@ -0,0 +1,125 @@
import { useEffect, useState } from "react";
import { api, clearToken, getToken } from "./api";
import Login from "./components/Login";
import ProductList from "./components/ProductList";
import ProductDetail from "./components/ProductDetail";
import SubmissionsPage from "./components/SubmissionsPage";
import { Inbox, LogOut, Package } from "lucide-react";
type Tab = "products" | "submissions";
type View = { name: "list" } | { name: "detail"; id: string };
export default function App() {
const [authed, setAuthed] = useState(false);
const [checking, setChecking] = useState(true);
const [username, setUsername] = useState("");
const [tab, setTab] = useState<Tab>("products");
const [pending, setPending] = useState<number | null>(null);
const [view, setView] = useState<View>({ name: "list" });
useEffect(() => {
if (!authed) return;
api
.listSubmissions("pending", 1, 1)
.then((r) => setPending(r.pending))
.catch(() => undefined);
}, [authed]);
useEffect(() => {
if (!getToken()) {
setChecking(false);
return;
}
api
.me()
.then((r) => {
setUsername(r.username);
setAuthed(true);
})
.catch(() => clearToken())
.finally(() => setChecking(false));
}, []);
function onLoggedIn(name: string) {
setUsername(name);
setAuthed(true);
setView({ name: "list" });
}
function logout() {
clearToken();
setAuthed(false);
setUsername("");
}
if (checking) {
return (
<div className="flex h-full items-center justify-center text-gray-500">
</div>
);
}
if (!authed) return <Login onLoggedIn={onLoggedIn} />;
return (
<div className="flex h-full flex-col">
<header className="flex items-center justify-between bg-white px-6 py-3 shadow-sm">
<div className="flex items-center gap-6">
<div className="flex items-center gap-2 text-lg font-semibold text-gray-800">
<Package className="h-5 w-5 text-emerald-600" />
·
</div>
<nav className="flex items-center gap-1 text-sm">
<button
onClick={() => {
setTab("products");
setView({ name: "list" });
}}
className={`px-3 py-1.5 rounded-md flex items-center gap-1.5 ${
tab === "products"
? "bg-emerald-50 text-emerald-700"
: "text-gray-600 hover:bg-gray-100"
}`}
>
<Package className="h-4 w-4" />
</button>
<button
onClick={() => setTab("submissions")}
className={`px-3 py-1.5 rounded-md flex items-center gap-1.5 ${
tab === "submissions"
? "bg-emerald-50 text-emerald-700"
: "text-gray-600 hover:bg-gray-100"
}`}
>
<Inbox className="h-4 w-4" /> 稿
{pending != null && pending > 0 && (
<span className="ml-1 text-xs bg-amber-500 text-white rounded-full px-1.5">
{pending}
</span>
)}
</button>
</nav>
</div>
<div className="flex items-center gap-4 text-sm text-gray-600">
<span>{username}</span>
<button
onClick={logout}
className="flex items-center gap-1 rounded px-2 py-1 text-gray-500 hover:bg-gray-100 hover:text-gray-800"
>
<LogOut className="h-4 w-4" /> 退
</button>
</div>
</header>
<main className="flex-1 overflow-auto p-6">
{tab === "submissions" ? (
<SubmissionsPage onPending={setPending} />
) : view.name === "list" ? (
<ProductList onOpen={(id) => setView({ name: "detail", id })} />
) : (
<ProductDetail id={view.id} onBack={() => setView({ name: "list" })} />
)}
</main>
</div>
);
}
+112
View File
@@ -0,0 +1,112 @@
// API base derives from Vite's BASE_URL (/ping/) so it matches the nginx prefix.
const API_BASE = `${import.meta.env.BASE_URL}api`;
const TOKEN_KEY = "opengoods_admin_token";
export function getToken(): string | null {
return localStorage.getItem(TOKEN_KEY);
}
export function setToken(token: string) {
localStorage.setItem(TOKEN_KEY, token);
}
export function clearToken() {
localStorage.removeItem(TOKEN_KEY);
}
export class ApiError extends Error {
status: number;
constructor(status: number, message: string) {
super(message);
this.status = status;
}
}
async function request<T>(path: string, options: RequestInit = {}): Promise<T> {
const headers: Record<string, string> = {
"Content-Type": "application/json",
...(options.headers as Record<string, string>),
};
const token = getToken();
if (token) headers.Authorization = `Bearer ${token}`;
const res = await fetch(`${API_BASE}${path}`, { ...options, headers });
if (res.status === 401) {
clearToken();
throw new ApiError(401, "登录已过期,请重新登录");
}
const text = await res.text();
const data = text ? JSON.parse(text) : null;
if (!res.ok) {
const msg = data?.error?.message || `请求失败 (${res.status})`;
throw new ApiError(res.status, msg);
}
return data as T;
}
export const api = {
login: (username: string, password: string) =>
request<{ token: string; username: string }>("/login", {
method: "POST",
body: JSON.stringify({ username, password }),
}),
me: () => request<{ username: string }>("/me"),
listProducts: (q: string, page: number, size: number) =>
request<{
items: import("./types").ProductRow[];
page: number;
size: number;
total: number;
completeness_fields: string[];
}>(`/products?q=${encodeURIComponent(q)}&page=${page}&size=${size}`),
getProduct: (id: string) =>
request<import("./types").ProductDetail>(`/products/${id}`),
updateProduct: (id: string, body: unknown) =>
request<import("./types").ProductDetail>(`/products/${id}`, {
method: "PUT",
body: JSON.stringify(body),
}),
listAudit: (id: string) =>
request<{ items: import("./types").AuditEntry[] }>(`/products/${id}/audit`),
addImage: (id: string, url: string, kind: string) =>
request<import("./types").ProductImage>(`/products/${id}/images`, {
method: "POST",
body: JSON.stringify({ url, kind }),
}),
deleteImage: (id: string, imageId: string) =>
request<{ status: string }>(`/products/${id}/images/${imageId}`, {
method: "DELETE",
}),
addMsrp: (id: string, body: unknown) =>
request<import("./types").MSRP>(`/products/${id}/msrp`, {
method: "POST",
body: JSON.stringify(body),
}),
deleteMsrp: (id: string, msrpId: string) =>
request<{ status: string }>(`/products/${id}/msrp/${msrpId}`, {
method: "DELETE",
}),
listBrands: () =>
request<{ items: import("./types").Brand[] }>("/brands"),
listCategories: () =>
request<{ items: import("./types").Category[] }>("/categories"),
listSubmissions: (status: string, page: number, size: number) =>
request<{
items: import("./types").SubmissionRow[];
page: number;
size: number;
total: number;
pending: number;
}>(`/submissions?status=${encodeURIComponent(status)}&page=${page}&size=${size}`),
getSubmission: (id: string) =>
request<import("./types").SubmissionDetail>(`/submissions/${id}`),
approveSubmission: (id: string) =>
request<import("./types").ProductDetail>(`/submissions/${id}/approve`, {
method: "POST",
}),
rejectSubmission: (id: string, note: string) =>
request<{ status: string }>(`/submissions/${id}/reject`, {
method: "POST",
body: JSON.stringify({ note }),
}),
};
+75
View File
@@ -0,0 +1,75 @@
import { useState } from "react";
import { api, setToken } from "../api";
import { Package } from "lucide-react";
export default function Login({
onLoggedIn,
}: {
onLoggedIn: (username: string) => void;
}) {
const [username, setUsername] = useState("");
const [password, setPassword] = useState("");
const [error, setError] = useState("");
const [loading, setLoading] = useState(false);
async function submit(e: React.FormEvent) {
e.preventDefault();
setError("");
setLoading(true);
try {
const r = await api.login(username, password);
setToken(r.token);
onLoggedIn(r.username);
} catch (err) {
setError(err instanceof Error ? err.message : "登录失败");
} finally {
setLoading(false);
}
}
return (
<div className="flex h-full items-center justify-center">
<form
onSubmit={submit}
className="w-80 rounded-xl bg-white p-8 shadow-md"
>
<div className="mb-6 flex flex-col items-center gap-2">
<Package className="h-8 w-8 text-emerald-600" />
<h1 className="text-lg font-semibold text-gray-800">
·
</h1>
</div>
{error && (
<div className="mb-4 rounded bg-red-50 px-3 py-2 text-sm text-red-600">
{error}
</div>
)}
<label className="mb-3 block">
<span className="mb-1 block text-sm text-gray-600"></span>
<input
value={username}
onChange={(e) => setUsername(e.target.value)}
className="w-full rounded border border-gray-300 px-3 py-2 text-sm focus:border-emerald-500 focus:outline-none"
autoFocus
/>
</label>
<label className="mb-5 block">
<span className="mb-1 block text-sm text-gray-600"></span>
<input
type="password"
value={password}
onChange={(e) => setPassword(e.target.value)}
className="w-full rounded border border-gray-300 px-3 py-2 text-sm focus:border-emerald-500 focus:outline-none"
/>
</label>
<button
type="submit"
disabled={loading}
className="w-full rounded bg-emerald-600 py-2 text-sm font-medium text-white hover:bg-emerald-700 disabled:opacity-60"
>
{loading ? "登录中…" : "登录"}
</button>
</form>
</div>
);
}
@@ -0,0 +1,642 @@
import { useEffect, useMemo, useState } from "react";
import { api } from "../api";
import {
AuditEntry,
Brand,
Category,
FIELD_LABELS,
ProductDetail as Detail,
} from "../types";
import {
ArrowLeft,
Plus,
Save,
Trash2,
AlertCircle,
History,
} from "lucide-react";
const NUTRIMENT_KEYS: { key: string; label: string }[] = [
{ key: "energy_kcal", label: "能量 (kcal)" },
{ key: "energy_kj", label: "能量 (kJ)" },
{ key: "fat", label: "脂肪 (g)" },
{ key: "saturated_fat", label: "饱和脂肪 (g)" },
{ key: "carbohydrates", label: "碳水 (g)" },
{ key: "sugars", label: "糖 (g)" },
{ key: "proteins", label: "蛋白质 (g)" },
{ key: "salt", label: "盐 (g)" },
];
const STATUS_OPTIONS = [
{ value: "active", label: "在用" },
{ value: "merged", label: "已合并" },
{ value: "deprecated", label: "已停用" },
];
const ACTION_LABEL: Record<string, string> = {
update: "编辑",
add_image: "新增图片",
delete_image: "删除图片",
add_msrp: "新增建议零售价",
delete_msrp: "删除建议零售价",
};
function Card({
title,
children,
}: {
title: string;
children: React.ReactNode;
}) {
return (
<div className="rounded-lg border border-gray-200 bg-white p-5">
<h3 className="mb-4 text-sm font-semibold text-gray-700">{title}</h3>
{children}
</div>
);
}
function Field({
label,
children,
}: {
label: string;
children: React.ReactNode;
}) {
return (
<label className="block">
<span className="mb-1 block text-xs text-gray-500">{label}</span>
{children}
</label>
);
}
const inputCls =
"w-full rounded border border-gray-300 px-3 py-2 text-sm focus:border-emerald-500 focus:outline-none";
export default function ProductDetail({
id,
onBack,
}: {
id: string;
onBack: () => void;
}) {
const [d, setD] = useState<Detail | null>(null);
const [brands, setBrands] = useState<Brand[]>([]);
const [categories, setCategories] = useState<Category[]>([]);
const [audit, setAudit] = useState<AuditEntry[]>([]);
const [loading, setLoading] = useState(true);
const [saving, setSaving] = useState(false);
const [msg, setMsg] = useState("");
const [error, setError] = useState("");
// editable form state
const [name, setName] = useState("");
const [gtin, setGtin] = useState("");
const [brandName, setBrandName] = useState("");
const [categoryId, setCategoryId] = useState("");
const [netValue, setNetValue] = useState("");
const [netUnit, setNetUnit] = useState("");
const [country, setCountry] = useState("");
const [status, setStatus] = useState("active");
const [ingredients, setIngredients] = useState("");
const [allergens, setAllergens] = useState("");
const [additives, setAdditives] = useState("");
const [nutriments, setNutriments] = useState<Record<string, string>>({});
const [basis, setBasis] = useState("");
const [serving, setServing] = useState("");
const [nutriScore, setNutriScore] = useState("");
function hydrate(detail: Detail) {
setD(detail);
setName(detail.name);
setGtin(detail.gtin || "");
setBrandName(detail.brand || "");
setCategoryId(detail.category_id || "");
setNetValue(detail.net_content_value?.toString() || "");
setNetUnit(detail.net_content_unit || "");
setCountry(detail.country_of_origin || "");
setStatus(detail.status);
setIngredients(detail.ingredients_text || "");
setAllergens(detail.allergens.join(", "));
setAdditives(detail.additives.join(", "));
const nm: Record<string, string> = {};
if (detail.nutriments) {
for (const [k, v] of Object.entries(detail.nutriments)) nm[k] = String(v);
}
setNutriments(nm);
setBasis(detail.nutrition_basis || "");
setServing(detail.serving_size || "");
setNutriScore(detail.nutri_score || "");
}
function reload() {
setLoading(true);
Promise.all([api.getProduct(id), api.listAudit(id)])
.then(([detail, a]) => {
hydrate(detail);
setAudit(a.items);
})
.catch((e) => setError(e.message))
.finally(() => setLoading(false));
}
useEffect(() => {
reload();
api.listBrands().then((r) => setBrands(r.items)).catch(() => {});
api.listCategories().then((r) => setCategories(r.items)).catch(() => {});
// eslint-disable-next-line react-hooks/exhaustive-deps
}, [id]);
const missing = useMemo(() => d?.missing ?? [], [d]);
function parseList(s: string): string[] {
return s
.split(",")
.map((x) => x.trim())
.filter(Boolean);
}
async function save() {
setSaving(true);
setMsg("");
setError("");
const nm: Record<string, number> = {};
for (const [k, v] of Object.entries(nutriments)) {
const n = parseFloat(v);
if (!Number.isNaN(n)) nm[k] = n;
}
const body = {
gtin: gtin.trim() || null,
name: name.trim(),
brand_name: brandName.trim() || null,
brand_id: brandName.trim() ? undefined : null,
category_id: categoryId || null,
net_content_value: netValue.trim() ? parseFloat(netValue) : null,
net_content_unit: netUnit.trim() || null,
country_of_origin: country.trim() || null,
status,
ingredients_text: ingredients.trim() || null,
allergens: parseList(allergens),
additives: parseList(additives),
nutriments: nm,
nutrition_basis: basis || null,
serving_size: serving.trim() || null,
nutri_score: nutriScore || null,
};
try {
const updated = await api.updateProduct(id, body);
hydrate(updated);
const a = await api.listAudit(id);
setAudit(a.items);
setMsg("已保存");
setTimeout(() => setMsg(""), 2500);
} catch (e) {
setError(e instanceof Error ? e.message : "保存失败");
} finally {
setSaving(false);
}
}
if (loading) {
return <div className="text-gray-400"></div>;
}
if (!d) {
return (
<div>
<button onClick={onBack} className="text-emerald-600">
</button>
<p className="mt-4 text-red-600">{error || "未找到商品"}</p>
</div>
);
}
return (
<div className="mx-auto max-w-5xl space-y-5">
<div className="flex items-center justify-between">
<button
onClick={onBack}
className="flex items-center gap-1 text-sm text-gray-600 hover:text-gray-900"
>
<ArrowLeft className="h-4 w-4" />
</button>
<div className="flex items-center gap-3">
{msg && <span className="text-sm text-emerald-600">{msg}</span>}
{error && <span className="text-sm text-red-600">{error}</span>}
<span className="text-xs text-gray-400">
{Math.round(d.quality_score * 100)}
</span>
<button
onClick={save}
disabled={saving}
className="flex items-center gap-1 rounded bg-emerald-600 px-4 py-2 text-sm text-white hover:bg-emerald-700 disabled:opacity-60"
>
<Save className="h-4 w-4" /> {saving ? "保存中…" : "保存"}
</button>
</div>
</div>
{missing.length > 0 && (
<div className="flex items-center gap-2 rounded-lg border border-amber-200 bg-amber-50 px-4 py-3 text-sm text-amber-700">
<AlertCircle className="h-4 w-4" />
{missing.map((f) => FIELD_LABELS[f] || f).join("、")}
</div>
)}
<Card title="基础信息">
<div className="grid grid-cols-2 gap-4">
<Field label="名称 *">
<input
className={inputCls}
value={name}
onChange={(e) => setName(e.target.value)}
/>
</Field>
<Field label="条码 (GTIN)">
<input
className={inputCls}
value={gtin}
onChange={(e) => setGtin(e.target.value)}
/>
</Field>
<Field label="品牌(不存在将自动创建)">
<input
className={inputCls}
list="brand-list"
value={brandName}
onChange={(e) => setBrandName(e.target.value)}
/>
<datalist id="brand-list">
{brands.map((b) => (
<option key={b.id} value={b.name} />
))}
</datalist>
</Field>
<Field label="品类">
<select
className={inputCls}
value={categoryId}
onChange={(e) => setCategoryId(e.target.value)}
>
<option value=""></option>
{categories.map((c) => (
<option key={c.id} value={c.id}>
{"\u00A0".repeat(c.level * 2)}
{c.name_zh} ({c.path})
</option>
))}
</select>
</Field>
<Field label="净含量">
<input
className={inputCls}
type="number"
step="any"
value={netValue}
onChange={(e) => setNetValue(e.target.value)}
/>
</Field>
<Field label="净含量单位 (g/ml/cl…)">
<input
className={inputCls}
value={netUnit}
onChange={(e) => setNetUnit(e.target.value)}
/>
</Field>
<Field label="产地">
<input
className={inputCls}
value={country}
onChange={(e) => setCountry(e.target.value)}
/>
</Field>
<Field label="状态">
<select
className={inputCls}
value={status}
onChange={(e) => setStatus(e.target.value)}
>
{STATUS_OPTIONS.map((o) => (
<option key={o.value} value={o.value}>
{o.label}
</option>
))}
</select>
</Field>
</div>
</Card>
<Card title="配料与营养">
<div className="mb-4 grid grid-cols-2 gap-4">
<Field label="配料表">
<textarea
className={inputCls}
rows={3}
value={ingredients}
onChange={(e) => setIngredients(e.target.value)}
/>
</Field>
<div className="grid grid-cols-2 gap-4">
<Field label="过敏原(逗号分隔)">
<input
className={inputCls}
value={allergens}
onChange={(e) => setAllergens(e.target.value)}
/>
</Field>
<Field label="添加剂(逗号分隔)">
<input
className={inputCls}
value={additives}
onChange={(e) => setAdditives(e.target.value)}
/>
</Field>
<Field label="营养基准">
<select
className={inputCls}
value={basis}
onChange={(e) => setBasis(e.target.value)}
>
<option value=""></option>
<option value="per_100g"> 100g</option>
<option value="per_100ml"> 100ml</option>
<option value="per_serving"></option>
</select>
</Field>
<Field label="份量">
<input
className={inputCls}
value={serving}
onChange={(e) => setServing(e.target.value)}
/>
</Field>
<Field label="Nutri-Score (A-E)">
<input
className={inputCls}
maxLength={1}
value={nutriScore}
onChange={(e) =>
setNutriScore(e.target.value.toUpperCase())
}
/>
</Field>
</div>
</div>
<div className="grid grid-cols-4 gap-3">
{NUTRIMENT_KEYS.map((n) => (
<Field key={n.key} label={n.label}>
<input
className={inputCls}
type="number"
step="any"
value={nutriments[n.key] ?? ""}
onChange={(e) =>
setNutriments((prev) => ({ ...prev, [n.key]: e.target.value }))
}
/>
</Field>
))}
</div>
</Card>
<ImagesCard
product={d}
onChange={reload}
onError={setError}
/>
<MsrpCard product={d} onChange={reload} onError={setError} />
<Card title="操作记录">
{audit.length === 0 ? (
<p className="text-sm text-gray-400"></p>
) : (
<ul className="space-y-2 text-sm">
{audit.map((a) => (
<li
key={a.id}
className="flex items-center gap-3 text-gray-600"
>
<History className="h-3.5 w-3.5 text-gray-400" />
<span className="text-gray-400">{a.created_at}</span>
<span className="font-medium text-gray-700">{a.actor}</span>
<span>{ACTION_LABEL[a.action] || a.action}</span>
{a.fields.length > 0 && (
<span className="text-gray-400">
[{a.fields.map((f) => FIELD_LABELS[f] || f).join("、")}]
</span>
)}
</li>
))}
</ul>
)}
</Card>
</div>
);
}
function ImagesCard({
product,
onChange,
onError,
}: {
product: Detail;
onChange: () => void;
onError: (m: string) => void;
}) {
const [url, setUrl] = useState("");
const [kind, setKind] = useState("front");
async function add() {
if (!url.trim()) return;
try {
await api.addImage(product.id, url.trim(), kind);
setUrl("");
onChange();
} catch (e) {
onError(e instanceof Error ? e.message : "添加失败");
}
}
async function remove(imageId: string) {
try {
await api.deleteImage(product.id, imageId);
onChange();
} catch (e) {
onError(e instanceof Error ? e.message : "删除失败");
}
}
return (
<Card title="图片(仅存 URL">
<div className="mb-3 flex flex-wrap gap-3">
{product.images.length === 0 && (
<span className="text-sm text-gray-400"></span>
)}
{product.images.map((im) => (
<div
key={im.id}
className="relative h-24 w-24 overflow-hidden rounded border border-gray-200"
>
<img
src={im.url}
alt={im.kind}
className="h-full w-full object-cover"
/>
<button
onClick={() => remove(im.id)}
className="absolute right-1 top-1 rounded bg-black/50 p-1 text-white hover:bg-black/70"
>
<Trash2 className="h-3 w-3" />
</button>
<span className="absolute bottom-0 left-0 bg-black/50 px-1 text-[10px] text-white">
{im.kind}
</span>
</div>
))}
</div>
<div className="flex items-center gap-2">
<input
className={inputCls}
placeholder="图片 URL"
value={url}
onChange={(e) => setUrl(e.target.value)}
/>
<select
className="rounded border border-gray-300 px-2 py-2 text-sm"
value={kind}
onChange={(e) => setKind(e.target.value)}
>
<option value="front"></option>
<option value="ingredients"></option>
<option value="nutrition"></option>
<option value="other"></option>
</select>
<button
onClick={add}
className="flex items-center gap-1 whitespace-nowrap rounded bg-gray-700 px-3 py-2 text-sm text-white hover:bg-gray-800"
>
<Plus className="h-4 w-4" />
</button>
</div>
</Card>
);
}
function MsrpCard({
product,
onChange,
onError,
}: {
product: Detail;
onChange: () => void;
onError: (m: string) => void;
}) {
const [amount, setAmount] = useState("");
const [currency, setCurrency] = useState("CNY");
const [region, setRegion] = useState("CN");
const [date, setDate] = useState("");
const [note, setNote] = useState("");
async function add() {
const a = parseFloat(amount);
if (Number.isNaN(a)) return;
try {
await api.addMsrp(product.id, {
amount: a,
currency,
region,
effective_date: date || null,
note: note.trim() || null,
});
setAmount("");
setNote("");
setDate("");
onChange();
} catch (e) {
onError(e instanceof Error ? e.message : "添加失败");
}
}
async function remove(msrpId: string) {
try {
await api.deleteMsrp(product.id, msrpId);
onChange();
} catch (e) {
onError(e instanceof Error ? e.message : "删除失败");
}
}
return (
<Card title="官方建议零售价(MSRP 快照,非售卖)">
<div className="mb-3 space-y-2">
{product.msrp.length === 0 && (
<span className="text-sm text-gray-400"></span>
)}
{product.msrp.map((m) => (
<div
key={m.id}
className="flex items-center gap-3 rounded border border-gray-100 bg-gray-50 px-3 py-2 text-sm"
>
<span className="font-medium text-gray-800">
{m.amount} {m.currency}
</span>
<span className="text-gray-500">{m.region}</span>
<span className="text-gray-400">{m.effective_date || ""}</span>
<span className="flex-1 text-gray-400">{m.note || ""}</span>
<button
onClick={() => remove(m.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="金额">
<input
className="w-28 rounded border border-gray-300 px-3 py-2 text-sm"
type="number"
step="any"
value={amount}
onChange={(e) => setAmount(e.target.value)}
/>
</Field>
<Field label="币种">
<input
className="w-20 rounded border border-gray-300 px-3 py-2 text-sm"
value={currency}
onChange={(e) => setCurrency(e.target.value.toUpperCase())}
/>
</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>
<Field label="生效日期">
<input
className="rounded border border-gray-300 px-3 py-2 text-sm"
type="date"
value={date}
onChange={(e) => setDate(e.target.value)}
/>
</Field>
<Field label="备注">
<input
className="w-40 rounded border border-gray-300 px-3 py-2 text-sm"
value={note}
onChange={(e) => setNote(e.target.value)}
/>
</Field>
<button
onClick={add}
className="flex items-center gap-1 rounded bg-gray-700 px-3 py-2 text-sm text-white hover:bg-gray-800"
>
<Plus className="h-4 w-4" />
</button>
</div>
</Card>
);
}
@@ -0,0 +1,178 @@
import { useEffect, useState } from "react";
import { api } from "../api";
import { FIELD_LABELS, ProductRow } from "../types";
import { Search, AlertCircle } from "lucide-react";
const STATUS_LABEL: Record<string, string> = {
active: "在用",
merged: "已合并",
deprecated: "已停用",
};
function QualityBadge({ score }: { score: number }) {
const pct = Math.round(score * 100);
const color =
score >= 0.8
? "bg-emerald-100 text-emerald-700"
: score >= 0.5
? "bg-amber-100 text-amber-700"
: "bg-red-100 text-red-700";
return (
<span className={`rounded px-2 py-0.5 text-xs font-medium ${color}`}>
{pct}
</span>
);
}
export default function ProductList({
onOpen,
}: {
onOpen: (id: string) => void;
}) {
const [q, setQ] = useState("");
const [input, setInput] = useState("");
const [page, setPage] = useState(1);
const [size] = useState(20);
const [rows, setRows] = useState<ProductRow[]>([]);
const [total, setTotal] = useState(0);
const [loading, setLoading] = useState(false);
const [error, setError] = useState("");
useEffect(() => {
setLoading(true);
setError("");
api
.listProducts(q, page, size)
.then((r) => {
setRows(r.items);
setTotal(r.total);
})
.catch((e) => setError(e.message))
.finally(() => setLoading(false));
}, [q, page, size]);
const pages = Math.max(1, Math.ceil(total / size));
return (
<div className="mx-auto max-w-6xl">
<div className="mb-4 flex items-center justify-between">
<h2 className="text-xl font-semibold text-gray-800">
<span className="text-sm font-normal text-gray-400"> {total} </span>
</h2>
<form
onSubmit={(e) => {
e.preventDefault();
setPage(1);
setQ(input.trim());
}}
className="flex items-center gap-2"
>
<div className="relative">
<Search className="absolute left-2 top-2.5 h-4 w-4 text-gray-400" />
<input
value={input}
onChange={(e) => setInput(e.target.value)}
placeholder="按名称 / 条码搜索"
className="w-64 rounded border border-gray-300 py-2 pl-8 pr-3 text-sm focus:border-emerald-500 focus:outline-none"
/>
</div>
<button className="rounded bg-emerald-600 px-3 py-2 text-sm text-white hover:bg-emerald-700">
</button>
</form>
</div>
{error && (
<div className="mb-3 rounded bg-red-50 px-3 py-2 text-sm text-red-600">
{error}
</div>
)}
<div className="overflow-hidden rounded-lg border border-gray-200 bg-white">
<table className="w-full text-sm">
<thead className="bg-gray-50 text-left text-xs uppercase text-gray-500">
<tr>
<th className="px-4 py-3"></th>
<th className="px-4 py-3"></th>
<th className="px-4 py-3"></th>
<th className="px-4 py-3"></th>
<th className="px-4 py-3"></th>
<th className="px-4 py-3"></th>
<th className="px-4 py-3"></th>
</tr>
</thead>
<tbody className="divide-y divide-gray-100">
{loading ? (
<tr>
<td colSpan={7} className="px-4 py-8 text-center text-gray-400">
</td>
</tr>
) : rows.length === 0 ? (
<tr>
<td colSpan={7} className="px-4 py-8 text-center text-gray-400">
</td>
</tr>
) : (
rows.map((r) => (
<tr
key={r.id}
onClick={() => onOpen(r.id)}
className="cursor-pointer hover:bg-emerald-50/50"
>
<td className="px-4 py-3 font-medium text-gray-800">{r.name}</td>
<td className="px-4 py-3 text-gray-600">{r.brand || "—"}</td>
<td className="px-4 py-3 font-mono text-xs text-gray-500">
{r.gtin || "—"}
</td>
<td className="px-4 py-3 text-xs text-gray-500">
{r.category_path || "—"}
</td>
<td className="px-4 py-3 text-gray-600">
{STATUS_LABEL[r.status] || r.status}
</td>
<td className="px-4 py-3">
<QualityBadge score={r.quality_score} />
</td>
<td className="px-4 py-3">
{r.missing.length === 0 ? (
<span className="text-xs text-emerald-600"></span>
) : (
<span className="flex items-center gap-1 text-xs text-amber-600">
<AlertCircle className="h-3.5 w-3.5" />
{r.missing
.map((f) => FIELD_LABELS[f] || f)
.join("、")}
</span>
)}
</td>
</tr>
))
)}
</tbody>
</table>
</div>
<div className="mt-4 flex items-center justify-end gap-2 text-sm text-gray-600">
<button
disabled={page <= 1}
onClick={() => setPage((p) => p - 1)}
className="rounded border border-gray-300 px-3 py-1 disabled:opacity-50"
>
</button>
<span>
{page} / {pages}
</span>
<button
disabled={page >= pages}
onClick={() => setPage((p) => p + 1)}
className="rounded border border-gray-300 px-3 py-1 disabled:opacity-50"
>
</button>
</div>
</div>
);
}
@@ -0,0 +1,314 @@
import { useEffect, useState } from "react";
import { api, ApiError } from "../api";
import type { SubmissionDetail, SubmissionRow } from "../types";
import { FIELD_LABELS } from "../types";
import { ArrowLeft, Check, X } from "lucide-react";
const STATUS_TABS = [
{ key: "pending", label: "待审核" },
{ key: "approved", label: "已通过" },
{ key: "rejected", label: "已驳回" },
];
const STATUS_BADGE: Record<string, string> = {
pending: "bg-amber-50 text-amber-700",
approved: "bg-emerald-50 text-emerald-700",
rejected: "bg-red-50 text-red-700",
};
const STATUS_TEXT: Record<string, string> = {
pending: "待审核",
approved: "已通过",
rejected: "已驳回",
};
export default function SubmissionsPage({ onPending }: { onPending?: (n: number) => void }) {
const [tab, setTab] = useState("pending");
const [rows, setRows] = useState<SubmissionRow[]>([]);
const [openId, setOpenId] = useState<string | null>(null);
const [error, setError] = useState("");
async function load() {
setError("");
try {
const res = await api.listSubmissions(tab, 1, 50);
setRows(res.items);
onPending?.(res.pending);
} catch (e) {
setError(e instanceof ApiError ? e.message : "加载失败");
}
}
useEffect(() => {
load();
// eslint-disable-next-line react-hooks/exhaustive-deps
}, [tab]);
if (openId) {
return (
<SubmissionView
id={openId}
onBack={() => {
setOpenId(null);
load();
}}
/>
);
}
return (
<div>
<div className="flex items-center gap-2 mb-4">
{STATUS_TABS.map((t) => (
<button
key={t.key}
onClick={() => setTab(t.key)}
className={`px-3 py-1.5 rounded-md text-sm ${
tab === t.key
? "bg-emerald-600 text-white"
: "bg-white border text-gray-600 hover:bg-gray-50"
}`}
>
{t.label}
</button>
))}
</div>
{error && (
<div className="mb-3 bg-red-50 text-red-700 text-sm rounded px-4 py-2">{error}</div>
)}
<div className="bg-white border rounded-lg overflow-hidden">
<table className="w-full text-sm">
<thead className="bg-gray-50 text-gray-500 text-left">
<tr>
<th className="px-4 py-2 font-medium"></th>
<th className="px-4 py-2 font-medium"></th>
<th className="px-4 py-2 font-medium">稿</th>
<th className="px-4 py-2 font-medium"></th>
<th className="px-4 py-2 font-medium"></th>
<th className="px-4 py-2 font-medium"></th>
</tr>
</thead>
<tbody className="divide-y">
{rows.length === 0 ? (
<tr>
<td colSpan={6} className="px-4 py-8 text-center text-gray-400">
稿
</td>
</tr>
) : (
rows.map((r) => (
<tr
key={r.id}
onClick={() => setOpenId(r.id)}
className="cursor-pointer hover:bg-gray-50"
>
<td className="px-4 py-2 text-gray-800">{r.name}</td>
<td className="px-4 py-2 text-gray-500">{r.gtin || "—"}</td>
<td className="px-4 py-2 text-gray-500">{r.submitter_name || "匿名"}</td>
<td className="px-4 py-2">
<span className="text-xs text-gray-500">
{r.matched ? "补全已有商品" : "新建商品"}
</span>
</td>
<td className="px-4 py-2 text-gray-500">
{new Date(r.created_at).toLocaleString()}
</td>
<td className="px-4 py-2">
<span className={`text-xs rounded px-2 py-0.5 ${STATUS_BADGE[r.status]}`}>
{STATUS_TEXT[r.status]}
</span>
</td>
</tr>
))
)}
</tbody>
</table>
</div>
</div>
);
}
function Field({ label, value }: { label: string; value: React.ReactNode }) {
if (value === null || value === undefined || value === "") return null;
return (
<div className="flex py-1.5 text-sm">
<div className="w-28 shrink-0 text-gray-400">{label}</div>
<div className="text-gray-800 break-all">{value}</div>
</div>
);
}
function SubmissionView({ id, onBack }: { id: string; onBack: () => void }) {
const [d, setD] = useState<SubmissionDetail | null>(null);
const [error, setError] = useState("");
const [busy, setBusy] = useState(false);
const [rejecting, setRejecting] = useState(false);
const [reason, setReason] = useState("");
useEffect(() => {
api.getSubmission(id).then(setD).catch((e) => setError(e.message));
}, [id]);
async function approve() {
if (!confirm("确认通过该投稿?将写入正式商品库并记录来源 community。")) return;
setBusy(true);
setError("");
try {
await api.approveSubmission(id);
onBack();
} catch (e) {
setError(e instanceof ApiError ? e.message : "操作失败");
} finally {
setBusy(false);
}
}
async function reject() {
setBusy(true);
setError("");
try {
await api.rejectSubmission(id, reason.trim());
onBack();
} catch (e) {
setError(e instanceof ApiError ? e.message : "操作失败");
} finally {
setBusy(false);
}
}
if (error && !d) {
return (
<div>
<button onClick={onBack} className="text-sm text-gray-500 flex items-center gap-1 mb-4">
<ArrowLeft className="h-4 w-4" />
</button>
<div className="bg-red-50 text-red-700 text-sm rounded px-4 py-3">{error}</div>
</div>
);
}
if (!d) return <div className="text-gray-400"></div>;
const p = d.payload;
const nutri = Object.entries(p.nutriments || {});
return (
<div className="max-w-3xl">
<button onClick={onBack} className="text-sm text-gray-500 flex items-center gap-1 mb-4">
<ArrowLeft className="h-4 w-4" /> 稿
</button>
<div className="flex items-center justify-between">
<h2 className="text-lg font-semibold text-gray-800">{p.name}</h2>
<span className={`text-xs rounded px-2 py-0.5 ${STATUS_BADGE[d.status]}`}>
{STATUS_TEXT[d.status]}
</span>
</div>
{error && <div className="mt-3 bg-red-50 text-red-700 text-sm rounded px-4 py-2">{error}</div>}
{d.target_product_id && (
<div className="mt-3 text-sm bg-blue-50 text-blue-700 rounded px-4 py-2">
<b></b>
</div>
)}
<div className="bg-white border rounded-lg p-5 mt-4">
<h3 className="font-medium text-gray-700 mb-2">稿</h3>
<Field label="商品名称" value={p.name} />
<Field label="条码" value={p.gtin} />
<Field label="品牌" value={p.brand_name} />
<Field
label="净含量"
value={p.net_content_value != null ? `${p.net_content_value} ${p.net_content_unit || ""}` : null}
/>
<Field label="产地" value={p.country_of_origin} />
<Field label="配料" value={p.ingredients_text} />
{nutri.length > 0 && (
<Field
label="营养成分"
value={
<span>
{p.nutrition_basis ? `(${p.nutrition_basis}) ` : ""}
{nutri.map(([k, v]) => `${FIELD_LABELS[k] || k}:${v}`).join("")}
</span>
}
/>
)}
{p.images && p.images.length > 0 && (
<Field
label="图片"
value={
<div className="flex flex-wrap gap-2">
{p.images.map((im, i) => (
<a key={i} href={im.url} target="_blank" rel="noreferrer">
<img src={im.url} alt="" className="h-20 w-20 object-cover rounded border" />
</a>
))}
</div>
}
/>
)}
</div>
<div className="bg-white border rounded-lg p-5 mt-4">
<h3 className="font-medium text-gray-700 mb-2">稿 / </h3>
<Field label="称呼" value={p.submitter_name || "匿名"} />
<Field label="联系方式" value={p.submitter_contact} />
<Field label="备注" value={p.note} />
<Field label="提交时间" value={new Date(d.created_at).toLocaleString()} />
{d.reviewed_by && <Field label="审核人" value={d.reviewed_by} />}
{d.review_note && <Field label="驳回原因" value={d.review_note} />}
{d.result_product_id && <Field label="收录商品ID" value={d.result_product_id} />}
</div>
{d.status === "pending" && (
<div className="mt-5">
{rejecting ? (
<div className="bg-white border rounded-lg p-4">
<label className="block text-xs text-gray-500 mb-1"></label>
<input
className="w-full border rounded-md px-3 py-2 text-sm"
value={reason}
onChange={(e) => setReason(e.target.value)}
placeholder="例如:资料无法核实 / 重复投稿"
/>
<div className="mt-3 flex gap-2">
<button
onClick={reject}
disabled={busy}
className="px-4 py-2 rounded bg-red-600 text-white text-sm hover:bg-red-700 disabled:opacity-60"
>
</button>
<button
onClick={() => setRejecting(false)}
className="px-4 py-2 rounded border text-sm text-gray-600"
>
</button>
</div>
</div>
) : (
<div className="flex gap-3">
<button
onClick={approve}
disabled={busy}
className="px-5 py-2.5 rounded-lg bg-emerald-600 text-white font-medium hover:bg-emerald-700 disabled:opacity-60 flex items-center gap-1.5"
>
<Check className="h-4 w-4" />
</button>
<button
onClick={() => setRejecting(true)}
disabled={busy}
className="px-5 py-2.5 rounded-lg border text-gray-700 font-medium hover:bg-gray-50 flex items-center gap-1.5"
>
<X className="h-4 w-4" />
</button>
</div>
)}
</div>
)}
</div>
);
}
+16
View File
@@ -0,0 +1,16 @@
@tailwind base;
@tailwind components;
@tailwind utilities;
html,
body,
#root {
height: 100%;
}
body {
margin: 0;
background: #f3f4f6;
font-family: system-ui, -apple-system, "Segoe UI", Roboto, "Helvetica Neue",
Arial, "PingFang SC", "Microsoft YaHei", sans-serif;
}
+10
View File
@@ -0,0 +1,10 @@
import React from "react";
import ReactDOM from "react-dom/client";
import App from "./App";
import "./index.css";
ReactDOM.createRoot(document.getElementById("root")!).render(
<React.StrictMode>
<App />
</React.StrictMode>,
);
+140
View File
@@ -0,0 +1,140 @@
export interface ProductRow {
id: string;
gtin: string | null;
name: string;
brand: string | null;
category_path: string | null;
status: string;
quality_score: number;
missing: string[];
updated_at: string;
}
export interface ProductImage {
id: string;
url: string;
kind: string;
license: string | null;
}
export interface MSRP {
id: string;
amount: number;
currency: string;
region: string;
effective_date: string | null;
source_url: string | null;
note: string | null;
}
export interface ProductDetail {
id: string;
gtin: string | null;
name: string;
brand_id: string | null;
brand: string | null;
category_id: string | null;
category_path: string | null;
net_content_value: number | null;
net_content_unit: string | null;
country_of_origin: string | null;
status: string;
quality_score: number;
ingredients_text: string | null;
allergens: string[];
additives: string[];
nutriments: Record<string, number> | null;
nutrition_basis: string | null;
serving_size: string | null;
nutri_score: string | null;
images: ProductImage[];
msrp: MSRP[];
missing: string[];
updated_at: string;
}
export interface Brand {
id: string;
name: string;
}
export interface Category {
id: string;
name_zh: string;
name_en: string | null;
path: string;
level: number;
}
export interface AuditEntry {
id: string;
actor: string;
action: string;
fields: string[];
created_at: string;
}
export interface SubmissionRow {
id: string;
gtin: string | null;
name: string;
status: string;
submitter_name: string | null;
matched: boolean;
created_at: string;
reviewed_at: string | null;
}
export interface SubmissionImage {
url: string;
kind: string;
}
export interface SubmissionPayload {
gtin: string | null;
name: string;
brand_name: string | null;
category_id: string | null;
net_content_value: number | null;
net_content_unit: string | null;
country_of_origin: string | null;
ingredients_text: string | null;
nutriments: Record<string, number> | null;
nutrition_basis: string | null;
serving_size: string | null;
nutri_score: string | null;
images: SubmissionImage[] | null;
submitter_name: string | null;
submitter_contact: string | null;
note: string | null;
}
export interface SubmissionDetail {
id: string;
status: string;
gtin: string | null;
name: string;
submitter_name: string | null;
submitter_contact: string | null;
note: string | null;
review_note: string | null;
reviewed_by: string | null;
reviewed_at: string | null;
created_at: string;
target_product_id: string | null;
result_product_id: string | null;
payload: SubmissionPayload;
existing_product?: ProductDetail;
}
export const FIELD_LABELS: Record<string, string> = {
name: "名称",
gtin: "条码",
brand: "品牌",
category: "品类",
net_content: "净含量",
country_of_origin: "产地",
nutriments: "营养成分",
ingredients: "配料",
image: "图片",
};
+1
View File
@@ -0,0 +1 @@
/// <reference types="vite/client" />
+6
View File
@@ -0,0 +1,6 @@
/** @type {import('tailwindcss').Config} */
export default {
content: ["./index.html", "./src/**/*.{ts,tsx}"],
theme: { extend: {} },
plugins: [],
};
+21
View File
@@ -0,0 +1,21 @@
{
"compilerOptions": {
"target": "ES2020",
"useDefineForClassFields": true,
"lib": ["ES2020", "DOM", "DOM.Iterable"],
"module": "ESNext",
"skipLibCheck": true,
"moduleResolution": "bundler",
"allowImportingTsExtensions": true,
"resolveJsonModule": true,
"isolatedModules": true,
"noEmit": true,
"jsx": "react-jsx",
"strict": true,
"noUnusedLocals": true,
"noUnusedParameters": true,
"noFallthroughCasesInSwitch": true
},
"include": ["src"],
"references": [{ "path": "./tsconfig.node.json" }]
}
+11
View File
@@ -0,0 +1,11 @@
{
"compilerOptions": {
"composite": true,
"skipLibCheck": true,
"module": "ESNext",
"moduleResolution": "bundler",
"allowSyntheticDefaultImports": true,
"strict": true
},
"include": ["vite.config.ts"]
}
+9
View File
@@ -0,0 +1,9 @@
import { defineConfig } from "vite";
import react from "@vitejs/plugin-react";
// Served under /ping by the admin Go binary; base must match the nginx prefix.
export default defineConfig({
base: "/ping/",
plugins: [react()],
build: { outDir: "dist", emptyOutDir: true },
});
+29
View File
@@ -0,0 +1,29 @@
# Admin console image: builds the SPA (node), embeds it into the Go admin
# binary, and ships a static scratch runtime. Build context is the repo root.
# Stage 1: build the admin SPA.
FROM node:22-alpine AS web
WORKDIR /web
ENV npm_config_registry=https://registry.npmmirror.com
COPY admin-frontend/package.json admin-frontend/package-lock.json* ./
RUN npm ci || npm install
COPY admin-frontend/ ./
RUN npm run build
# Stage 2: build the Go admin binary with the SPA embedded.
FROM golang:1.23-alpine AS build
ENV GOPROXY=https://goproxy.cn,direct
WORKDIR /src
COPY api/go.mod api/go.sum ./
RUN go mod download
COPY api/ ./
RUN rm -rf internal/adminweb/dist && mkdir -p internal/adminweb/dist
COPY --from=web /web/dist/ internal/adminweb/dist/
RUN CGO_ENABLED=0 go build -o /out/admin ./cmd/admin
# Stage 3: minimal runtime.
FROM scratch
COPY --from=build /etc/ssl/certs/ca-certificates.crt /etc/ssl/certs/
COPY --from=build /out/admin /admin
EXPOSE 8080
ENTRYPOINT ["/admin"]
+31
View File
@@ -0,0 +1,31 @@
# Public API image: builds the public SPA (homepage + search + contribute),
# embeds it into the Go read-only server binary, and ships a static scratch
# runtime (used where gcr.io/distroless is not reachable). Build context is the
# repo root.
# Stage 1: build the public SPA.
FROM node:22-alpine AS web
WORKDIR /web
ENV npm_config_registry=https://registry.npmmirror.com
COPY public-frontend/package.json public-frontend/package-lock.json* ./
RUN npm ci || npm install
COPY public-frontend/ ./
RUN npm run build
# Stage 2: build the Go server binary with the SPA embedded.
FROM golang:1.23-alpine AS build
ENV GOPROXY=https://goproxy.cn,direct
WORKDIR /src
COPY api/go.mod api/go.sum ./
RUN go mod download
COPY api/ ./
RUN rm -rf internal/publicweb/dist && mkdir -p internal/publicweb/dist
COPY --from=web /web/dist/ internal/publicweb/dist/
RUN CGO_ENABLED=0 go build -o /out/server ./cmd/server
# Stage 3: minimal runtime.
FROM scratch
COPY --from=build /etc/ssl/certs/ca-certificates.crt /etc/ssl/certs/
COPY --from=build /out/server /server
EXPOSE 8080
ENTRYPOINT ["/server"]
+83
View File
@@ -0,0 +1,83 @@
// Command admin starts the OpenGoods admin console (authenticated write API +
// embedded SPA), served under a base path (default /ping).
package main
import (
"context"
"crypto/rand"
"log"
"net/http"
"os"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"golang.org/x/crypto/bcrypt"
"github.com/baicai2026-baicai/goods/api/internal/adminhandler"
"github.com/baicai2026-baicai/goods/api/internal/adminstore"
"github.com/baicai2026-baicai/goods/api/internal/adminweb"
"github.com/baicai2026-baicai/goods/api/internal/auth"
)
func getenv(key, fallback string) string {
if v, ok := os.LookupEnv(key); ok && v != "" {
return v
}
return fallback
}
func main() {
addr := getenv("GOODS_ADMIN_ADDR", ":8080")
dbURL := getenv("OPENGOODS_DATABASE_URL", "postgres://opengoods:opengoods@localhost:5432/opengoods?sslmode=disable")
basePath := getenv("GOODS_ADMIN_BASE_PATH", "/ping")
username := getenv("GOODS_ADMIN_USER", "admin")
// Password: prefer a bcrypt hash; otherwise hash a plaintext password.
var passwordHash []byte
if h := os.Getenv("GOODS_ADMIN_PASSWORD_HASH"); h != "" {
passwordHash = []byte(h)
} else if p := os.Getenv("GOODS_ADMIN_PASSWORD"); p != "" {
hashed, err := bcrypt.GenerateFromPassword([]byte(p), bcrypt.DefaultCost)
if err != nil {
log.Fatalf("failed to hash admin password: %v", err)
}
passwordHash = hashed
} else {
log.Fatal("set GOODS_ADMIN_PASSWORD or GOODS_ADMIN_PASSWORD_HASH")
}
secret := []byte(os.Getenv("GOODS_ADMIN_JWT_SECRET"))
if len(secret) == 0 {
secret = make([]byte, 32)
if _, err := rand.Read(secret); err != nil {
log.Fatalf("failed to generate jwt secret: %v", err)
}
log.Print("warning: GOODS_ADMIN_JWT_SECRET not set; using a random secret (tokens invalidate on restart)")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dbURL)
if err != nil {
log.Fatalf("failed to create db pool: %v", err)
}
defer pool.Close()
pingCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()
if err := pool.Ping(pingCtx); err != nil {
log.Printf("warning: database not reachable at startup: %v", err)
}
authn := auth.New(username, passwordHash, secret, 12*time.Hour)
h := adminhandler.New(adminstore.New(pool), authn, basePath, adminweb.Dist())
srv := &http.Server{
Addr: addr,
Handler: h.Router(),
ReadHeaderTimeout: 10 * time.Second,
}
log.Printf("OpenGoods admin console listening on %s (base path %s)", addr, basePath)
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
log.Fatalf("server error: %v", err)
}
}
+21 -1
View File
@@ -2,20 +2,40 @@
package main
import (
"context"
"log"
"net/http"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/baicai2026-baicai/goods/api/internal/config"
"github.com/baicai2026-baicai/goods/api/internal/handler"
"github.com/baicai2026-baicai/goods/api/internal/publicweb"
"github.com/baicai2026-baicai/goods/api/internal/store"
)
func main() {
cfg := config.Load()
ctx := context.Background()
pool, err := pgxpool.New(ctx, cfg.DatabaseURL)
if err != nil {
log.Fatalf("failed to create db pool: %v", err)
}
defer pool.Close()
pingCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()
if err := pool.Ping(pingCtx); err != nil {
log.Printf("warning: database not reachable at startup: %v", err)
}
h := handler.New(store.New(pool), publicweb.Dist())
srv := &http.Server{
Addr: cfg.Addr,
Handler: handler.Router(),
Handler: h.Router(),
ReadHeaderTimeout: 10 * time.Second,
}
+13 -1
View File
@@ -2,4 +2,16 @@ module github.com/baicai2026-baicai/goods/api
go 1.23.4
require github.com/go-chi/chi/v5 v5.1.0
require (
github.com/go-chi/chi/v5 v5.1.0
github.com/jackc/pgx/v5 v5.7.2
golang.org/x/crypto v0.31.0
)
require (
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
github.com/jackc/puddle/v2 v2.2.2 // indirect
golang.org/x/sync v0.10.0 // indirect
golang.org/x/text v0.21.0 // indirect
)
+28
View File
@@ -1,2 +1,30 @@
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/go-chi/chi/v5 v5.1.0 h1:acVI1TYaD+hhedDJ3r54HyA6sExp3HfXq7QWEEY/xMw=
github.com/go-chi/chi/v5 v5.1.0/go.mod h1:DslCQbL2OYiznFReuXYUmQ2hGd1aDpCnlMNITLSKoi8=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM=
github.com/jackc/pgx/v5 v5.7.2 h1:mLoDLV6sonKlvjIEsV56SkWNCnuNv531l94GaIzO+XI=
github.com/jackc/pgx/v5 v5.7.2/go.mod h1:ncY89UGWxg82EykZUwSpUKEfccBGGYq1xjrOpsbsfGQ=
github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo=
github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk=
github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4=
golang.org/x/crypto v0.31.0 h1:ihbySMvVjLAeSH1IbfcRTkD/iNscyz8rGzjF/E5hV6U=
golang.org/x/crypto v0.31.0/go.mod h1:kDsLvtWBEx7MV9tJOj9bnXsPbxwJQ6csT/x4KIN4Ssk=
golang.org/x/sync v0.10.0 h1:3NQrjDixjgGwUOCaF8w2+VYHv0Ve/vGYSbdkTa98gmQ=
golang.org/x/sync v0.10.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
golang.org/x/text v0.21.0 h1:zyQAAkrwaneQ066sspRyJaG9VNi/YJ1NfzcGB3hZ/qo=
golang.org/x/text v0.21.0/go.mod h1:4IBbMaMmOPCJ8SecivzSH54+73PCFmPWxNTLm+vZkEQ=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+390
View File
@@ -0,0 +1,390 @@
// Package adminhandler wires up the authenticated admin console: a JSON write
// API mounted under a base path (default /ping) plus the embedded SPA.
package adminhandler
import (
"encoding/json"
"errors"
"io/fs"
"net/http"
"strings"
"time"
"github.com/go-chi/chi/v5"
"github.com/go-chi/chi/v5/middleware"
"github.com/baicai2026-baicai/goods/api/internal/adminstore"
"github.com/baicai2026-baicai/goods/api/internal/auth"
)
// Handler holds the admin dependencies.
type Handler struct {
store *adminstore.Store
authn *auth.Authenticator
basePath string
spa fs.FS
submitLimit *rateLimiter
}
// New constructs an admin Handler. basePath is e.g. "/ping" (no trailing slash).
func New(store *adminstore.Store, authn *auth.Authenticator, basePath string, spa fs.FS) *Handler {
basePath = "/" + strings.Trim(basePath, "/")
return &Handler{
store: store,
authn: authn,
basePath: basePath,
spa: spa,
submitLimit: newRateLimiter(5, 10*time.Minute),
}
}
// Router builds the HTTP handler.
func (h *Handler) Router() http.Handler {
r := chi.NewRouter()
r.Use(middleware.RequestID)
r.Use(middleware.RealIP)
r.Use(middleware.Recoverer)
r.Route(h.basePath, func(r chi.Router) {
r.Get("/healthz", func(w http.ResponseWriter, _ *http.Request) {
writeJSON(w, http.StatusOK, map[string]string{"status": "ok"})
})
r.Post("/api/login", h.Login)
r.Group(func(r chi.Router) {
r.Use(h.authn.Middleware)
r.Get("/api/me", h.Me)
r.Get("/api/products", h.ListProducts)
r.Get("/api/products/{id}", h.GetProduct)
r.Put("/api/products/{id}", h.UpdateProduct)
r.Get("/api/products/{id}/audit", h.ListAudit)
r.Post("/api/products/{id}/images", h.AddImage)
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.Get("/api/brands", h.ListBrands)
r.Get("/api/categories", h.ListCategories)
r.Get("/api/submissions", h.ListSubmissions)
r.Get("/api/submissions/{id}", h.GetSubmission)
r.Post("/api/submissions/{id}/approve", h.ApproveSubmission)
r.Post("/api/submissions/{id}/reject", h.RejectSubmission)
})
r.Handle("/*", http.HandlerFunc(h.serveSPA))
})
// Public, unauthenticated contribution endpoint (proxied at /api/public/*).
// Submissions enter a moderation queue and never touch products until an
// admin approves them.
r.Post("/api/public/submissions", h.CreateSubmission)
return r
}
func (h *Handler) serveSPA(w http.ResponseWriter, r *http.Request) {
rel := strings.TrimPrefix(r.URL.Path, h.basePath)
rel = strings.TrimPrefix(rel, "/")
if rel == "" {
rel = "index.html"
}
if f, err := h.spa.Open(rel); err == nil {
f.Close()
http.StripPrefix(h.basePath+"/", http.FileServer(http.FS(h.spa))).ServeHTTP(w, r)
return
}
// SPA fallback: serve index.html for client-side routes.
index, err := h.spa.Open("index.html")
if err != nil {
http.NotFound(w, r)
return
}
defer index.Close()
data, _ := fs.ReadFile(h.spa, "index.html")
w.Header().Set("Content-Type", "text/html; charset=utf-8")
_, _ = w.Write(data)
}
// ---------- auth ----------
// Login authenticates and returns a bearer token.
func (h *Handler) Login(w http.ResponseWriter, r *http.Request) {
var body struct {
Username string `json:"username"`
Password string `json:"password"`
}
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
writeError(w, http.StatusBadRequest, "bad_request", "invalid body")
return
}
token, err := h.authn.Login(body.Username, body.Password)
if err != nil {
writeError(w, http.StatusUnauthorized, "unauthorized", "用户名或密码错误")
return
}
writeJSON(w, http.StatusOK, map[string]string{"token": token, "username": body.Username})
}
// Me returns the current authenticated user.
func (h *Handler) Me(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, map[string]string{"username": auth.UserFrom(r.Context())})
}
// ---------- products ----------
// ListProducts returns a paginated product list.
func (h *Handler) ListProducts(w http.ResponseWriter, r *http.Request) {
q := r.URL.Query().Get("q")
page, size := pageParams(r)
items, total, err := h.store.ListProducts(r.Context(), q, size, (page-1)*size)
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, map[string]any{
"items": items, "page": page, "size": size, "total": total,
"completeness_fields": adminstore.CompletenessFields,
})
}
// GetProduct returns full editable detail.
func (h *Handler) GetProduct(w http.ResponseWriter, r *http.Request) {
d, err := h.store.GetProduct(r.Context(), chi.URLParam(r, "id"))
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, d)
}
// UpdateProduct applies an edit.
func (h *Handler) UpdateProduct(w http.ResponseWriter, r *http.Request) {
var in adminstore.ProductInput
if err := json.NewDecoder(r.Body).Decode(&in); err != nil {
writeError(w, http.StatusBadRequest, "bad_request", "invalid body")
return
}
if strings.TrimSpace(in.Name) == "" {
writeError(w, http.StatusBadRequest, "bad_request", "名称不能为空")
return
}
d, err := h.store.UpdateProduct(r.Context(), chi.URLParam(r, "id"), auth.UserFrom(r.Context()), in)
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, d)
}
// ListAudit returns audit history for a product.
func (h *Handler) ListAudit(w http.ResponseWriter, r *http.Request) {
items, err := h.store.ListAudit(r.Context(), chi.URLParam(r, "id"), 100)
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, map[string]any{"items": items})
}
// AddImage adds an image URL.
func (h *Handler) AddImage(w http.ResponseWriter, r *http.Request) {
var body struct {
URL string `json:"url"`
Kind string `json:"kind"`
}
if err := json.NewDecoder(r.Body).Decode(&body); err != nil || strings.TrimSpace(body.URL) == "" {
writeError(w, http.StatusBadRequest, "bad_request", "图片 URL 不能为空")
return
}
im, err := h.store.AddImage(r.Context(), chi.URLParam(r, "id"), auth.UserFrom(r.Context()), body.URL, body.Kind)
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusCreated, im)
}
// DeleteImage removes an image.
func (h *Handler) DeleteImage(w http.ResponseWriter, r *http.Request) {
err := h.store.DeleteImage(r.Context(), chi.URLParam(r, "id"), chi.URLParam(r, "imageID"), auth.UserFrom(r.Context()))
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, map[string]string{"status": "deleted"})
}
// AddMSRP adds a suggested-retail-price snapshot.
func (h *Handler) AddMSRP(w http.ResponseWriter, r *http.Request) {
var in adminstore.MSRPInput
if err := json.NewDecoder(r.Body).Decode(&in); err != nil {
writeError(w, http.StatusBadRequest, "bad_request", "invalid body")
return
}
m, err := h.store.AddMSRP(r.Context(), chi.URLParam(r, "id"), auth.UserFrom(r.Context()), in)
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusCreated, m)
}
// DeleteMSRP removes an MSRP snapshot.
func (h *Handler) DeleteMSRP(w http.ResponseWriter, r *http.Request) {
err := h.store.DeleteMSRP(r.Context(), chi.URLParam(r, "id"), chi.URLParam(r, "msrpID"), auth.UserFrom(r.Context()))
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, map[string]string{"status": "deleted"})
}
// ListBrands returns brand options.
func (h *Handler) ListBrands(w http.ResponseWriter, r *http.Request) {
items, err := h.store.ListBrands(r.Context())
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, map[string]any{"items": items})
}
// ListCategories returns category options.
func (h *Handler) ListCategories(w http.ResponseWriter, r *http.Request) {
items, err := h.store.ListCategories(r.Context())
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, map[string]any{"items": items})
}
// ---------- submissions ----------
// CreateSubmission accepts an anonymous public contribution into the queue.
func (h *Handler) CreateSubmission(w http.ResponseWriter, r *http.Request) {
if !h.submitLimit.allow(realIP(r)) {
writeError(w, http.StatusTooManyRequests, "rate_limited", "提交过于频繁,请稍后再试")
return
}
var in adminstore.SubmissionInput
if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, 1<<20)).Decode(&in); err != nil {
writeError(w, http.StatusBadRequest, "bad_request", "invalid body")
return
}
if strings.TrimSpace(in.Name) == "" {
writeError(w, http.StatusBadRequest, "bad_request", "商品名称不能为空")
return
}
id, err := h.store.CreateSubmission(r.Context(), in, realIP(r))
if err != nil {
writeError(w, http.StatusInternalServerError, "internal_error", err.Error())
return
}
writeJSON(w, http.StatusCreated, map[string]string{"id": id, "status": "pending"})
}
// ListSubmissions returns the moderation queue (admin).
func (h *Handler) ListSubmissions(w http.ResponseWriter, r *http.Request) {
status := r.URL.Query().Get("status")
page, size := pageParams(r)
items, total, err := h.store.ListSubmissions(r.Context(), status, size, (page-1)*size)
if h.handleErr(w, err) {
return
}
pending, _ := h.store.PendingSubmissionCount(r.Context())
writeJSON(w, http.StatusOK, map[string]any{
"items": items, "page": page, "size": size, "total": total, "pending": pending,
})
}
// GetSubmission returns full submission detail (admin).
func (h *Handler) GetSubmission(w http.ResponseWriter, r *http.Request) {
d, err := h.store.GetSubmission(r.Context(), chi.URLParam(r, "id"))
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, d)
}
// ApproveSubmission applies a contribution to the product store (admin).
func (h *Handler) ApproveSubmission(w http.ResponseWriter, r *http.Request) {
d, err := h.store.ApproveSubmission(r.Context(), chi.URLParam(r, "id"), auth.UserFrom(r.Context()))
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, d)
}
// RejectSubmission rejects a contribution with a reviewer note (admin).
func (h *Handler) RejectSubmission(w http.ResponseWriter, r *http.Request) {
var body struct {
Note string `json:"note"`
}
_ = json.NewDecoder(r.Body).Decode(&body)
err := h.store.RejectSubmission(r.Context(), chi.URLParam(r, "id"), auth.UserFrom(r.Context()), body.Note)
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, map[string]string{"status": "rejected"})
}
// ---------- helpers ----------
func realIP(r *http.Request) string {
if ip := r.Header.Get("X-Forwarded-For"); ip != "" {
if i := strings.IndexByte(ip, ','); i >= 0 {
return strings.TrimSpace(ip[:i])
}
return strings.TrimSpace(ip)
}
if ip := r.Header.Get("X-Real-IP"); ip != "" {
return ip
}
return r.RemoteAddr
}
func (h *Handler) handleErr(w http.ResponseWriter, err error) bool {
if err == nil {
return false
}
if errors.Is(err, adminstore.ErrNotFound) {
writeError(w, http.StatusNotFound, "not_found", "资源不存在")
return true
}
if errors.Is(err, adminstore.ErrConflict) {
writeError(w, http.StatusConflict, "conflict", "该投稿已被处理")
return true
}
writeError(w, http.StatusInternalServerError, "internal_error", err.Error())
return true
}
func pageParams(r *http.Request) (page, size int) {
page = atoiDefault(r.URL.Query().Get("page"), 1)
if page < 1 {
page = 1
}
size = atoiDefault(r.URL.Query().Get("size"), 20)
if size < 1 {
size = 20
}
if size > 100 {
size = 100
}
return page, size
}
func atoiDefault(s string, fallback int) int {
if s == "" {
return fallback
}
n := 0
for _, c := range s {
if c < '0' || c > '9' {
return fallback
}
n = n*10 + int(c-'0')
}
return n
}
func writeJSON(w http.ResponseWriter, status int, body any) {
w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(body)
}
func writeError(w http.ResponseWriter, status int, code, message string) {
writeJSON(w, status, map[string]any{"error": map[string]string{"code": code, "message": message}})
}
+41
View File
@@ -0,0 +1,41 @@
package adminhandler
import (
"sync"
"time"
)
// rateLimiter is a simple fixed-window per-key limiter used to throttle
// anonymous public submissions (basic anti-spam; captcha can be added later).
type rateLimiter struct {
mu sync.Mutex
hits map[string][]time.Time
limit int
window time.Duration
}
func newRateLimiter(limit int, window time.Duration) *rateLimiter {
return &rateLimiter{hits: map[string][]time.Time{}, limit: limit, window: window}
}
// allow reports whether the key may proceed, recording the hit if so.
func (r *rateLimiter) allow(key string) bool {
now := time.Now()
cutoff := now.Add(-r.window)
r.mu.Lock()
defer r.mu.Unlock()
kept := r.hits[key][:0]
for _, t := range r.hits[key] {
if t.After(cutoff) {
kept = append(kept, t)
}
}
if len(kept) >= r.limit {
r.hits[key] = kept
return false
}
r.hits[key] = append(kept, now)
return true
}
+282
View File
@@ -0,0 +1,282 @@
// Package adminstore is the read/write data-access layer for the admin console.
// Unlike the public store (read-only), it performs INSERT/UPDATE/DELETE and
// records field-level provenance (source = "manual") plus an audit_log entry
// for every write, then recomputes product.quality_score.
package adminstore
import (
"context"
"encoding/json"
"errors"
"strconv"
"time"
"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 pgx pool for admin operations.
type Store struct {
pool *pgxpool.Pool
}
// New constructs an admin Store.
func New(pool *pgxpool.Pool) *Store { return &Store{pool: pool} }
// Ping verifies DB connectivity.
func (s *Store) Ping(ctx context.Context) error { return s.pool.Ping(ctx) }
// CompletenessFields mirrors ingestion/opengoods/etl/quality.py COMPLETENESS_FIELDS.
var CompletenessFields = []string{
"name", "gtin", "brand", "category", "net_content",
"country_of_origin", "nutriments", "ingredients", "image",
}
// ---------- list ----------
// ProductRow is a list-view row for the admin product table.
type ProductRow struct {
ID string `json:"id"`
GTIN *string `json:"gtin"`
Name string `json:"name"`
Brand *string `json:"brand"`
CategoryPath *string `json:"category_path"`
Status string `json:"status"`
QualityScore float64 `json:"quality_score"`
Missing []string `json:"missing"`
UpdatedAt string `json:"updated_at"`
}
// ListProducts returns a paginated, optionally name/gtin-filtered list.
func (s *Store) ListProducts(ctx context.Context, q string, limit, offset int) ([]ProductRow, int, error) {
args := []any{}
where := "WHERE 1=1"
if q != "" {
args = append(args, q)
where += " AND (p.name ILIKE '%' || $1 || '%' OR p.gtin ILIKE '%' || $1 || '%')"
}
var total int
if err := s.pool.QueryRow(ctx, "SELECT count(*) FROM product p "+where, args...).Scan(&total); err != nil {
return nil, 0, err
}
args = append(args, limit, offset)
sql := `
SELECT p.id, p.gtin, p.name, b.name, c.path::text, p.status, p.quality_score,
p.updated_at,
(p.brand_id IS NOT NULL) AS has_brand,
(p.category_id IS NOT NULL) AS has_cat,
(p.net_content_canonical IS NOT NULL) AS has_net,
(p.country_of_origin IS NOT NULL AND p.country_of_origin <> '') AS has_country,
(f.nutriments IS NOT NULL AND f.nutriments::text <> '{}') AS has_nutri,
(f.ingredients_text IS NOT NULL AND f.ingredients_text <> '') AS has_ing,
EXISTS (SELECT 1 FROM product_image pi WHERE pi.product_id = p.id) AS has_img
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 ` + where +
" ORDER BY p.updated_at DESC LIMIT $" + strconv.Itoa(len(args)-1) + " OFFSET $" + strconv.Itoa(len(args))
rows, err := s.pool.Query(ctx, sql, args...)
if err != nil {
return nil, 0, err
}
defer rows.Close()
out := []ProductRow{}
for rows.Next() {
var r ProductRow
var hasBrand, hasCat, hasNet, hasCountry, hasNutri, hasIng, hasImg bool
var updated time.Time
if err := rows.Scan(&r.ID, &r.GTIN, &r.Name, &r.Brand, &r.CategoryPath, &r.Status,
&r.QualityScore, &updated, &hasBrand, &hasCat, &hasNet, &hasCountry,
&hasNutri, &hasIng, &hasImg); err != nil {
return nil, 0, err
}
r.UpdatedAt = updated.Format(time.RFC3339)
present := map[string]bool{
"name": r.Name != "",
"gtin": r.GTIN != nil && *r.GTIN != "",
"brand": hasBrand,
"category": hasCat,
"net_content": hasNet,
"country_of_origin": hasCountry,
"nutriments": hasNutri,
"ingredients": hasIng,
"image": hasImg,
}
r.Missing = []string{}
for _, f := range CompletenessFields {
if !present[f] {
r.Missing = append(r.Missing, f)
}
}
out = append(out, r)
}
return out, total, rows.Err()
}
// ---------- detail ----------
// ProductImage is one image row.
type ProductImage struct {
ID string `json:"id"`
URL string `json:"url"`
Kind string `json:"kind"`
License *string `json:"license"`
}
// MSRP is one suggested-retail-price snapshot.
type MSRP struct {
ID string `json:"id"`
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"`
}
// ProductDetail is the full editable view of a product.
type ProductDetail struct {
ID string `json:"id"`
GTIN *string `json:"gtin"`
Name string `json:"name"`
BrandID *string `json:"brand_id"`
Brand *string `json:"brand"`
CategoryID *string `json:"category_id"`
CategoryPath *string `json:"category_path"`
NetContentValue *float64 `json:"net_content_value"`
NetContentUnit *string `json:"net_content_unit"`
CountryOfOrigin *string `json:"country_of_origin"`
Status string `json:"status"`
QualityScore float64 `json:"quality_score"`
IngredientsText *string `json:"ingredients_text"`
Allergens []string `json:"allergens"`
Additives []string `json:"additives"`
Nutriments map[string]any `json:"nutriments"`
NutritionBasis *string `json:"nutrition_basis"`
ServingSize *string `json:"serving_size"`
NutriScore *string `json:"nutri_score"`
Images []ProductImage `json:"images"`
MSRP []MSRP `json:"msrp"`
Missing []string `json:"missing"`
UpdatedAt string `json:"updated_at"`
}
// GetProduct returns the full editable detail for one product.
func (s *Store) GetProduct(ctx context.Context, id string) (*ProductDetail, error) {
var d ProductDetail
var nutriments []byte
var updated time.Time
err := s.pool.QueryRow(ctx, `
SELECT p.id, p.gtin, p.name, p.brand_id, b.name, p.category_id, c.path::text,
p.net_content_value, p.net_content_unit, p.country_of_origin, p.status,
p.quality_score, p.updated_at,
f.ingredients_text, f.allergens, f.additives, f.nutriments,
f.nutrition_basis, f.serving_size, f.nutri_score
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
WHERE p.id = $1`, id).Scan(
&d.ID, &d.GTIN, &d.Name, &d.BrandID, &d.Brand, &d.CategoryID, &d.CategoryPath,
&d.NetContentValue, &d.NetContentUnit, &d.CountryOfOrigin, &d.Status,
&d.QualityScore, &updated,
&d.IngredientsText, &d.Allergens, &d.Additives, &nutriments,
&d.NutritionBasis, &d.ServingSize, &d.NutriScore,
)
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrNotFound
}
if err != nil {
return nil, err
}
d.UpdatedAt = updated.Format(time.RFC3339)
if len(nutriments) > 0 {
_ = json.Unmarshal(nutriments, &d.Nutriments)
}
if d.Allergens == nil {
d.Allergens = []string{}
}
if d.Additives == nil {
d.Additives = []string{}
}
imgs, err := s.listImages(ctx, id)
if err != nil {
return nil, err
}
d.Images = imgs
msrps, err := s.listMSRP(ctx, id)
if err != nil {
return nil, err
}
d.MSRP = msrps
d.Missing = missingFromDetail(&d)
return &d, nil
}
func missingFromDetail(d *ProductDetail) []string {
present := map[string]bool{
"name": d.Name != "",
"gtin": d.GTIN != nil && *d.GTIN != "",
"brand": d.BrandID != nil,
"category": d.CategoryID != nil,
"net_content": d.NetContentValue != nil,
"country_of_origin": d.CountryOfOrigin != nil && *d.CountryOfOrigin != "",
"nutriments": len(d.Nutriments) > 0,
"ingredients": d.IngredientsText != nil && *d.IngredientsText != "",
"image": len(d.Images) > 0,
}
missing := []string{}
for _, f := range CompletenessFields {
if !present[f] {
missing = append(missing, f)
}
}
return missing
}
func (s *Store) listImages(ctx context.Context, productID string) ([]ProductImage, error) {
rows, err := s.pool.Query(ctx,
"SELECT id, url, kind, license FROM product_image WHERE product_id = $1 ORDER BY id", productID)
if err != nil {
return nil, err
}
defer rows.Close()
out := []ProductImage{}
for rows.Next() {
var im ProductImage
if err := rows.Scan(&im.ID, &im.URL, &im.Kind, &im.License); err != nil {
return nil, err
}
out = append(out, im)
}
return out, rows.Err()
}
func (s *Store) listMSRP(ctx context.Context, productID string) ([]MSRP, error) {
rows, err := s.pool.Query(ctx, `
SELECT id, amount, currency, region, effective_date::text, source_url, note
FROM product_msrp WHERE product_id = $1 ORDER BY effective_date DESC NULLS LAST`, productID)
if err != nil {
return nil, err
}
defer rows.Close()
out := []MSRP{}
for rows.Next() {
var m MSRP
if err := rows.Scan(&m.ID, &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()
}
+137
View File
@@ -0,0 +1,137 @@
package adminstore
import (
"context"
"math"
"time"
"github.com/jackc/pgx/v5"
)
// Quality weights mirror ingestion/opengoods/etl/quality.py.
const (
wCompleteness = 0.4
wSourceTrust = 0.3
wAgreement = 0.2
wFreshness = 0.1
)
// queryer is satisfied by both *pgxpool.Pool and pgx.Tx.
type queryer interface {
QueryRow(ctx context.Context, sql string, args ...any) pgx.Row
}
func agreementFromSources(n int) float64 {
switch {
case n <= 1:
return 0.5
case n == 2:
return 0.8
default:
return 1.0
}
}
func freshnessFromAge(ageDays *float64) float64 {
if ageDays == nil {
return 0.5
}
d := *ageDays
switch {
case d <= 30:
return 1.0
case d <= 180:
return 0.8
case d <= 365:
return 0.6
case d <= 730:
return 0.4
default:
return 0.2
}
}
func (s *Store) computeQuality(ctx context.Context, q queryer, productID string) (float64, error) {
var name, country *string
var gtin *string
var brandID, categoryID *string
var netCanonical *float64
var ingredients *string
var hasNutri, hasImage bool
err := q.QueryRow(ctx, `
SELECT p.name, p.gtin, p.brand_id, p.category_id, p.net_content_canonical,
p.country_of_origin, f.ingredients_text,
(f.nutriments IS NOT NULL AND f.nutriments::text <> '{}'),
EXISTS (SELECT 1 FROM product_image pi WHERE pi.product_id = p.id)
FROM product p LEFT JOIN food_detail f ON f.product_id = p.id
WHERE p.id = $1`, productID).Scan(
&name, &gtin, &brandID, &categoryID, &netCanonical, &country,
&ingredients, &hasNutri, &hasImage)
if err != nil {
return 0, err
}
present := 0
bump := func(ok bool) {
if ok {
present++
}
}
bump(name != nil && *name != "")
bump(gtin != nil && *gtin != "")
bump(brandID != nil)
bump(categoryID != nil)
bump(netCanonical != nil)
bump(country != nil && *country != "")
bump(hasNutri)
bump(ingredients != nil && *ingredients != "")
bump(hasImage)
completeness := float64(present) / float64(len(CompletenessFields))
var sourceCount int
var sourceTrust *float64
var lastFetched *time.Time
err = q.QueryRow(ctx, `
SELECT count(DISTINCT ps.source_id), COALESCE(max(s.trust_weight),0), max(ps.fetched_at)
FROM product_source ps LEFT JOIN source s ON s.id = ps.source_id
WHERE ps.product_id = $1`, productID).Scan(&sourceCount, &sourceTrust, &lastFetched)
if err != nil {
return 0, err
}
trust := 0.0
if sourceTrust != nil {
trust = *sourceTrust
}
var ageDays *float64
if lastFetched != nil {
d := time.Since(*lastFetched).Hours() / 24.0
if d < 0 {
d = 0
}
ageDays = &d
}
raw := wCompleteness*completeness + wSourceTrust*trust +
wAgreement*agreementFromSources(sourceCount) + wFreshness*freshnessFromAge(ageDays)
raw = math.Max(0, math.Min(1, raw))
return math.Round(raw*1000) / 1000, nil
}
func (s *Store) recomputeQualityTx(ctx context.Context, tx pgx.Tx, productID string) (float64, error) {
v, err := s.computeQuality(ctx, tx, productID)
if err != nil {
return 0, err
}
_, err = tx.Exec(ctx, "UPDATE product SET quality_score=$1 WHERE id=$2", v, productID)
return v, err
}
func (s *Store) recomputeQuality(ctx context.Context, productID string) (float64, error) {
v, err := s.computeQuality(ctx, s.pool, productID)
if err != nil {
return 0, err
}
_, err = s.pool.Exec(ctx, "UPDATE product SET quality_score=$1 WHERE id=$2", v, productID)
return v, err
}
+30
View File
@@ -0,0 +1,30 @@
package adminstore
import "testing"
func TestAgreementFromSources(t *testing.T) {
cases := map[int]float64{0: 0.5, 1: 0.5, 2: 0.8, 3: 1.0, 9: 1.0}
for n, want := range cases {
if got := agreementFromSources(n); got != want {
t.Errorf("agreementFromSources(%d) = %v, want %v", n, got, want)
}
}
}
func TestFreshnessFromAge(t *testing.T) {
mk := func(d float64) *float64 { return &d }
if got := freshnessFromAge(nil); got != 0.5 {
t.Errorf("nil age = %v, want 0.5", got)
}
cases := []struct {
days float64
want float64
}{
{10, 1.0}, {30, 1.0}, {100, 0.8}, {300, 0.6}, {500, 0.4}, {1000, 0.2},
}
for _, c := range cases {
if got := freshnessFromAge(mk(c.days)); got != c.want {
t.Errorf("freshnessFromAge(%v) = %v, want %v", c.days, got, c.want)
}
}
}
+421
View File
@@ -0,0 +1,421 @@
package adminstore
import (
"context"
"encoding/json"
"errors"
"strconv"
"strings"
"time"
"github.com/jackc/pgx/v5"
)
// ErrConflict is returned when a submission has already been reviewed.
var ErrConflict = errors.New("conflict")
// SubmissionImage is one proposed image URL inside a contribution.
type SubmissionImage struct {
URL string `json:"url"`
Kind string `json:"kind"`
}
// SubmissionInput is the public contribution payload (no login required).
type SubmissionInput struct {
GTIN *string `json:"gtin"`
Name string `json:"name"`
BrandName *string `json:"brand_name"`
CategoryID *string `json:"category_id"`
NetContentValue *float64 `json:"net_content_value"`
NetContentUnit *string `json:"net_content_unit"`
CountryOfOrigin *string `json:"country_of_origin"`
IngredientsText *string `json:"ingredients_text"`
Nutriments map[string]any `json:"nutriments"`
NutritionBasis *string `json:"nutrition_basis"`
ServingSize *string `json:"serving_size"`
NutriScore *string `json:"nutri_score"`
Images []SubmissionImage `json:"images"`
MSRP []MSRPInput `json:"msrp"`
SubmitterName *string `json:"submitter_name"`
SubmitterContact *string `json:"submitter_contact"`
Note *string `json:"note"`
}
// SubmissionRow is a queue-list row for the admin review table.
type SubmissionRow struct {
ID string `json:"id"`
GTIN *string `json:"gtin"`
Name string `json:"name"`
Status string `json:"status"`
SubmitterName *string `json:"submitter_name"`
Matched bool `json:"matched"`
CreatedAt string `json:"created_at"`
ReviewedAt *string `json:"reviewed_at"`
}
// SubmissionDetail is the full review view of one contribution.
type SubmissionDetail struct {
ID string `json:"id"`
Status string `json:"status"`
GTIN *string `json:"gtin"`
Name string `json:"name"`
SubmitterName *string `json:"submitter_name"`
SubmitterContact *string `json:"submitter_contact"`
Note *string `json:"note"`
ReviewNote *string `json:"review_note"`
ReviewedBy *string `json:"reviewed_by"`
ReviewedAt *string `json:"reviewed_at"`
CreatedAt string `json:"created_at"`
TargetProductID *string `json:"target_product_id"`
ResultProductID *string `json:"result_product_id"`
Payload SubmissionInput `json:"payload"`
ExistingProduct *ProductDetail `json:"existing_product,omitempty"`
}
// CreateSubmission validates and stores a public contribution as pending.
func (s *Store) CreateSubmission(ctx context.Context, in SubmissionInput, remoteIP string) (string, error) {
in.Name = strings.TrimSpace(in.Name)
if in.Name == "" {
return "", errors.New("商品名称不能为空")
}
if in.GTIN != nil {
g := strings.TrimSpace(*in.GTIN)
if g == "" {
in.GTIN = nil
} else {
in.GTIN = &g
}
}
// Link to an existing product when the barcode already exists (supplement).
var target *string
if in.GTIN != nil {
var pid string
err := s.pool.QueryRow(ctx, "SELECT id FROM product WHERE gtin = $1", *in.GTIN).Scan(&pid)
if err == nil {
target = &pid
} else if !errors.Is(err, pgx.ErrNoRows) {
return "", err
}
}
payload, err := json.Marshal(in)
if err != nil {
return "", err
}
var id string
err = s.pool.QueryRow(ctx, `
INSERT INTO submission (gtin, name, payload, target_product_id, submitter_name, submitter_contact, note, remote_ip)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8) RETURNING id`,
in.GTIN, in.Name, payload, target, in.SubmitterName, in.SubmitterContact, in.Note, remoteIP).Scan(&id)
return id, err
}
// ListSubmissions returns submissions filtered by status (empty = all).
func (s *Store) ListSubmissions(ctx context.Context, status string, limit, offset int) ([]SubmissionRow, int, error) {
args := []any{}
where := "WHERE 1=1"
if status != "" {
args = append(args, status)
where += " AND status = $1"
}
var total int
if err := s.pool.QueryRow(ctx, "SELECT count(*) FROM submission "+where, args...).Scan(&total); err != nil {
return nil, 0, err
}
args = append(args, limit, offset)
sql := `
SELECT id, gtin, name, status, submitter_name, (target_product_id IS NOT NULL),
created_at, reviewed_at
FROM submission ` + where +
" ORDER BY (status='pending') DESC, created_at DESC LIMIT $" +
strconv.Itoa(len(args)-1) + " OFFSET $" + strconv.Itoa(len(args))
rows, err := s.pool.Query(ctx, sql, args...)
if err != nil {
return nil, 0, err
}
defer rows.Close()
out := []SubmissionRow{}
for rows.Next() {
var r SubmissionRow
var created time.Time
var reviewed *time.Time
if err := rows.Scan(&r.ID, &r.GTIN, &r.Name, &r.Status, &r.SubmitterName, &r.Matched, &created, &reviewed); err != nil {
return nil, 0, err
}
r.CreatedAt = created.Format(time.RFC3339)
if reviewed != nil {
t := reviewed.Format(time.RFC3339)
r.ReviewedAt = &t
}
out = append(out, r)
}
return out, total, rows.Err()
}
// PendingSubmissionCount returns the number of submissions awaiting review.
func (s *Store) PendingSubmissionCount(ctx context.Context) (int, error) {
var n int
err := s.pool.QueryRow(ctx, "SELECT count(*) FROM submission WHERE status='pending'").Scan(&n)
return n, err
}
// GetSubmission returns the full review detail for one submission.
func (s *Store) GetSubmission(ctx context.Context, id string) (*SubmissionDetail, error) {
var d SubmissionDetail
var payload []byte
var created time.Time
var reviewed *time.Time
err := s.pool.QueryRow(ctx, `
SELECT id, status, gtin, name, submitter_name, submitter_contact, note,
review_note, reviewed_by, reviewed_at, created_at, target_product_id, result_product_id, payload
FROM submission WHERE id = $1`, id).Scan(
&d.ID, &d.Status, &d.GTIN, &d.Name, &d.SubmitterName, &d.SubmitterContact, &d.Note,
&d.ReviewNote, &d.ReviewedBy, &reviewed, &created, &d.TargetProductID, &d.ResultProductID, &payload,
)
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrNotFound
}
if err != nil {
return nil, err
}
d.CreatedAt = created.Format(time.RFC3339)
if reviewed != nil {
t := reviewed.Format(time.RFC3339)
d.ReviewedAt = &t
}
if len(payload) > 0 {
_ = json.Unmarshal(payload, &d.Payload)
}
if d.TargetProductID != nil {
if ep, err := s.GetProduct(ctx, *d.TargetProductID); err == nil {
d.ExistingProduct = ep
}
}
return &d, nil
}
// RejectSubmission marks a pending submission as rejected with a reviewer note.
func (s *Store) RejectSubmission(ctx context.Context, id, actor, note string) error {
ct, err := s.pool.Exec(ctx, `
UPDATE submission SET status='rejected', review_note=$2, reviewed_by=$3, reviewed_at=now()
WHERE id=$1 AND status='pending'`, id, note, actor)
if err != nil {
return err
}
if ct.RowsAffected() == 0 {
// Distinguish missing vs already-reviewed.
var st string
if e := s.pool.QueryRow(ctx, "SELECT status FROM submission WHERE id=$1", id).Scan(&st); errors.Is(e, pgx.ErrNoRows) {
return ErrNotFound
}
return ErrConflict
}
_ = s.writeAudit(ctx, actor, "reject_submission", "submission", &id, []string{}, nil, map[string]string{"review_note": note})
return nil
}
// ApproveSubmission applies a pending contribution to the product store
// (creating or supplementing a product), records community provenance + audit,
// recomputes quality, and marks the submission approved.
func (s *Store) ApproveSubmission(ctx context.Context, id, actor string) (*ProductDetail, error) {
sub, err := s.GetSubmission(ctx, id)
if err != nil {
return nil, err
}
if sub.Status != "pending" {
return nil, ErrConflict
}
in := sub.Payload
tx, err := s.pool.Begin(ctx)
if err != nil {
return nil, err
}
defer tx.Rollback(ctx)
communityID, err := s.sourceIDTx(ctx, tx, "community")
if err != nil {
return nil, err
}
// Resolve the target product (existing supplement vs new create).
productID := ""
if sub.TargetProductID != nil {
productID = *sub.TargetProductID
} else if in.GTIN != nil {
var pid string
if e := tx.QueryRow(ctx, "SELECT id FROM product WHERE gtin=$1", *in.GTIN).Scan(&pid); e == nil {
productID = pid
} else if !errors.Is(e, pgx.ErrNoRows) {
return nil, e
}
}
var brandID *string
if in.BrandName != nil && strings.TrimSpace(*in.BrandName) != "" {
bid, err := s.ensureBrand(ctx, tx, strings.TrimSpace(*in.BrandName))
if err != nil {
return nil, err
}
brandID = &bid
}
var gpc *string
if in.CategoryID != nil && *in.CategoryID != "" {
if err := tx.QueryRow(ctx, "SELECT gpc_brick_code FROM category WHERE id=$1", *in.CategoryID).Scan(&gpc); err != nil && !errors.Is(err, pgx.ErrNoRows) {
return nil, err
}
}
canonical, err := s.netCanonical(ctx, tx, in.NetContentValue, in.NetContentUnit)
if err != nil {
return nil, err
}
fields := submissionFields(in)
if productID == "" {
// Create a new product from the contribution.
err = tx.QueryRow(ctx, `
INSERT INTO product (gtin, name, brand_id, category_id, gpc_brick_code,
net_content_value, net_content_unit, net_content_canonical, country_of_origin, status)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,'active') RETURNING id`,
in.GTIN, in.Name, brandID, in.CategoryID, gpc,
in.NetContentValue, in.NetContentUnit, canonical, in.CountryOfOrigin).Scan(&productID)
if err != nil {
return nil, err
}
} else {
// Supplement an existing product: only overwrite fields the
// contribution actually provides (COALESCE keeps current values).
_, err = tx.Exec(ctx, `
UPDATE product SET
name=COALESCE(NULLIF($2,''), name),
brand_id=COALESCE($3, brand_id),
category_id=COALESCE($4, category_id),
gpc_brick_code=COALESCE($5, gpc_brick_code),
net_content_value=COALESCE($6, net_content_value),
net_content_unit=COALESCE($7, net_content_unit),
net_content_canonical=COALESCE($8, net_content_canonical),
country_of_origin=COALESCE($9, country_of_origin),
gtin=COALESCE($10, gtin)
WHERE id=$1`,
productID, in.Name, brandID, in.CategoryID, gpc,
in.NetContentValue, in.NetContentUnit, canonical, in.CountryOfOrigin, in.GTIN)
if err != nil {
return nil, err
}
}
// food_detail: upsert, preserving existing values where not provided.
var nutriJSON []byte
if len(in.Nutriments) > 0 {
nutriJSON, _ = json.Marshal(in.Nutriments)
}
_, err = tx.Exec(ctx, `
INSERT INTO food_detail (product_id, ingredients_text, nutriments, nutrition_basis, serving_size, nutri_score)
VALUES ($1,$2,$3,$4,$5,$6)
ON CONFLICT (product_id) DO UPDATE SET
ingredients_text=COALESCE(EXCLUDED.ingredients_text, food_detail.ingredients_text),
nutriments=COALESCE(EXCLUDED.nutriments, food_detail.nutriments),
nutrition_basis=COALESCE(EXCLUDED.nutrition_basis, food_detail.nutrition_basis),
serving_size=COALESCE(EXCLUDED.serving_size, food_detail.serving_size),
nutri_score=COALESCE(EXCLUDED.nutri_score, food_detail.nutri_score)`,
productID, in.IngredientsText, nutriJSON, in.NutritionBasis, in.ServingSize, in.NutriScore)
if err != nil {
return nil, err
}
for _, im := range in.Images {
url := strings.TrimSpace(im.URL)
if url == "" {
continue
}
kind := im.Kind
if kind != "front" && kind != "ingredients" && kind != "nutrition" {
kind = "other"
}
if _, err := tx.Exec(ctx, `
INSERT INTO product_image (product_id, url, kind, source_id) VALUES ($1,$2,$3,$4)`,
productID, url, kind, communityID); err != nil {
return nil, err
}
}
for _, m := range in.MSRP {
if m.Amount <= 0 {
continue
}
cur := m.Currency
if cur == "" {
cur = "CNY"
}
region := m.Region
if region == "" {
region = "CN"
}
if _, err := tx.Exec(ctx, `
INSERT INTO product_msrp (product_id, amount, currency, region, source_id, source_url, effective_date, note)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8)`,
productID, m.Amount, cur, region, communityID, m.SourceURL, m.EffectiveDate, m.Note); err != nil {
return nil, err
}
}
if _, err := s.recomputeQualityTx(ctx, tx, productID); err != nil {
return nil, err
}
if _, err := tx.Exec(ctx, `
UPDATE submission SET status='approved', reviewed_by=$2, reviewed_at=now(), result_product_id=$3
WHERE id=$1`, id, actor, productID); err != nil {
return nil, err
}
// Field-level provenance for the contributed fields (community source).
if len(fields) > 0 {
if _, err := tx.Exec(ctx, `
INSERT INTO product_source (product_id, source_id, url, fields, fetched_at, raw)
VALUES ($1,$2,NULL,$3,now(),NULL)`, productID, communityID, fields); err != nil {
return nil, err
}
}
if err := tx.Commit(ctx); err != nil {
return nil, err
}
_ = s.writeAudit(ctx, actor, "approve_submission", "product", &productID, fields,
map[string]string{"submission_id": id}, map[string]string{"product_id": productID})
return s.GetProduct(ctx, productID)
}
func (s *Store) sourceIDTx(ctx context.Context, tx pgx.Tx, name string) (string, error) {
var id string
err := tx.QueryRow(ctx, "SELECT id FROM source WHERE name=$1", name).Scan(&id)
return id, err
}
// submissionFields lists the product fields a contribution provides values for.
func submissionFields(in SubmissionInput) []string {
fields := []string{"name"}
add := func(name string, present bool) {
if present {
fields = append(fields, name)
}
}
add("gtin", in.GTIN != nil && *in.GTIN != "")
add("brand", in.BrandName != nil && strings.TrimSpace(*in.BrandName) != "")
add("category", in.CategoryID != nil && *in.CategoryID != "")
add("net_content", in.NetContentValue != nil)
add("country_of_origin", in.CountryOfOrigin != nil && *in.CountryOfOrigin != "")
add("ingredients", in.IngredientsText != nil && *in.IngredientsText != "")
add("nutriments", len(in.Nutriments) > 0)
add("image", len(in.Images) > 0)
return fields
}
+415
View File
@@ -0,0 +1,415 @@
package adminstore
import (
"context"
"encoding/json"
"errors"
"strings"
"github.com/jackc/pgx/v5"
)
// ProductInput is the editable payload accepted from the admin UI.
type ProductInput struct {
GTIN *string `json:"gtin"`
Name string `json:"name"`
BrandID *string `json:"brand_id"`
BrandName *string `json:"brand_name"`
CategoryID *string `json:"category_id"`
NetContentValue *float64 `json:"net_content_value"`
NetContentUnit *string `json:"net_content_unit"`
CountryOfOrigin *string `json:"country_of_origin"`
Status string `json:"status"`
IngredientsText *string `json:"ingredients_text"`
Allergens []string `json:"allergens"`
Additives []string `json:"additives"`
Nutriments map[string]any `json:"nutriments"`
NutritionBasis *string `json:"nutrition_basis"`
ServingSize *string `json:"serving_size"`
NutriScore *string `json:"nutri_score"`
}
func normBrand(name string) string { return strings.Join(strings.Fields(strings.ToLower(name)), " ") }
func (s *Store) ensureBrand(ctx context.Context, tx pgx.Tx, name string) (string, error) {
var id string
err := tx.QueryRow(ctx, `
INSERT INTO brand (name, normalized_name) VALUES ($1, $2)
ON CONFLICT (normalized_name) DO UPDATE SET name = brand.name
RETURNING id`, name, normBrand(name)).Scan(&id)
return id, err
}
func (s *Store) manualSourceID(ctx context.Context, tx pgx.Tx) (string, error) {
var id string
err := tx.QueryRow(ctx, "SELECT id FROM source WHERE name = 'manual'").Scan(&id)
return id, err
}
// netCanonical converts value+unit to the canonical base unit via the unit table.
func (s *Store) netCanonical(ctx context.Context, tx pgx.Tx, value *float64, unit *string) (*float64, error) {
if value == nil || unit == nil || *unit == "" {
return nil, nil
}
var factor *float64
err := tx.QueryRow(ctx, "SELECT to_canonical_factor FROM unit WHERE code = $1", *unit).Scan(&factor)
if errors.Is(err, pgx.ErrNoRows) || factor == nil {
return nil, nil
}
if err != nil {
return nil, err
}
c := *value * *factor
return &c, nil
}
// UpdateProduct applies an edit, records provenance + audit, and recomputes quality.
func (s *Store) UpdateProduct(ctx context.Context, id, actor string, in ProductInput) (*ProductDetail, error) {
before, err := s.GetProduct(ctx, id)
if err != nil {
return nil, err
}
tx, err := s.pool.Begin(ctx)
if err != nil {
return nil, err
}
defer tx.Rollback(ctx)
// Resolve brand (create-by-name takes precedence over id).
brandID := in.BrandID
if in.BrandName != nil && strings.TrimSpace(*in.BrandName) != "" {
bid, err := s.ensureBrand(ctx, tx, strings.TrimSpace(*in.BrandName))
if err != nil {
return nil, err
}
brandID = &bid
}
// Resolve category gpc brick code.
var gpc *string
if in.CategoryID != nil && *in.CategoryID != "" {
if err := tx.QueryRow(ctx, "SELECT gpc_brick_code FROM category WHERE id = $1", *in.CategoryID).Scan(&gpc); err != nil && !errors.Is(err, pgx.ErrNoRows) {
return nil, err
}
}
canonical, err := s.netCanonical(ctx, tx, in.NetContentValue, in.NetContentUnit)
if err != nil {
return nil, err
}
status := in.Status
if status == "" {
status = before.Status
}
_, err = tx.Exec(ctx, `
UPDATE product SET gtin=$1, name=$2, brand_id=$3, category_id=$4, gpc_brick_code=$5,
net_content_value=$6, net_content_unit=$7, net_content_canonical=$8,
country_of_origin=$9, status=$10
WHERE id=$11`,
in.GTIN, in.Name, brandID, in.CategoryID, gpc,
in.NetContentValue, in.NetContentUnit, canonical,
in.CountryOfOrigin, status, id)
if err != nil {
return nil, err
}
var nutriJSON []byte
if in.Nutriments != nil {
nutriJSON, _ = json.Marshal(in.Nutriments)
}
allergens := in.Allergens
if allergens == nil {
allergens = []string{}
}
additives := in.Additives
if additives == nil {
additives = []string{}
}
_, err = tx.Exec(ctx, `
INSERT INTO food_detail (product_id, ingredients_text, allergens, additives,
nutriments, nutrition_basis, serving_size, nutri_score)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8)
ON CONFLICT (product_id) DO UPDATE SET
ingredients_text=EXCLUDED.ingredients_text,
allergens=EXCLUDED.allergens,
additives=EXCLUDED.additives,
nutriments=EXCLUDED.nutriments,
nutrition_basis=EXCLUDED.nutrition_basis,
serving_size=EXCLUDED.serving_size,
nutri_score=EXCLUDED.nutri_score`,
id, in.IngredientsText, allergens, additives,
nutriJSON, in.NutritionBasis, in.ServingSize, in.NutriScore)
if err != nil {
return nil, err
}
if _, err := s.recomputeQualityTx(ctx, tx, id); err != nil {
return nil, err
}
if err := tx.Commit(ctx); err != nil {
return nil, err
}
after, err := s.GetProduct(ctx, id)
if err != nil {
return nil, err
}
changed := diffFields(before, after)
if len(changed) > 0 {
if err := s.recordProvenance(ctx, id, changed); err != nil {
return nil, err
}
}
if err := s.writeAudit(ctx, actor, "update", "product", &id, changed, before, after); err != nil {
return nil, err
}
return after, nil
}
func strEq(a, b *string) bool {
if a == nil && b == nil {
return true
}
if a == nil || b == nil {
return false
}
return *a == *b
}
func floatEq(a, b *float64) bool {
if a == nil && b == nil {
return true
}
if a == nil || b == nil {
return false
}
return *a == *b
}
func diffFields(a, b *ProductDetail) []string {
changed := []string{}
add := func(name string, eq bool) {
if !eq {
changed = append(changed, name)
}
}
add("gtin", strEq(a.GTIN, b.GTIN))
add("name", a.Name == b.Name)
add("brand", strEq(a.BrandID, b.BrandID))
add("category", strEq(a.CategoryID, b.CategoryID))
add("net_content", floatEq(a.NetContentValue, b.NetContentValue) && strEq(a.NetContentUnit, b.NetContentUnit))
add("country_of_origin", strEq(a.CountryOfOrigin, b.CountryOfOrigin))
add("status", a.Status == b.Status)
add("ingredients", strEq(a.IngredientsText, b.IngredientsText))
ja, _ := json.Marshal(a.Nutriments)
jb, _ := json.Marshal(b.Nutriments)
add("nutriments", string(ja) == string(jb))
add("nutrition_basis", strEq(a.NutritionBasis, b.NutritionBasis))
add("serving_size", strEq(a.ServingSize, b.ServingSize))
add("nutri_score", strEq(a.NutriScore, b.NutriScore))
return changed
}
func (s *Store) recordProvenance(ctx context.Context, productID string, fields []string) error {
var srcID string
if err := s.pool.QueryRow(ctx, "SELECT id FROM source WHERE name = 'manual'").Scan(&srcID); err != nil {
return err
}
_, err := s.pool.Exec(ctx, `
INSERT INTO product_source (product_id, source_id, url, fields, fetched_at, raw)
VALUES ($1, $2, NULL, $3, now(), NULL)`, productID, srcID, fields)
return err
}
func (s *Store) writeAudit(ctx context.Context, actor, action, entity string, entityID *string, fields []string, before, after any) error {
bj, _ := json.Marshal(before)
aj, _ := json.Marshal(after)
if fields == nil {
fields = []string{}
}
_, err := s.pool.Exec(ctx, `
INSERT INTO audit_log (actor, action, entity, entity_id, fields, before, after)
VALUES ($1,$2,$3,$4,$5,$6,$7)`, actor, action, entity, entityID, fields, bj, aj)
return err
}
// ---------- images ----------
// AddImage inserts an image URL (manual source) and recomputes quality.
func (s *Store) AddImage(ctx context.Context, productID, actor, url, kind string) (*ProductImage, error) {
if kind == "" {
kind = "other"
}
var srcID string
if err := s.pool.QueryRow(ctx, "SELECT id FROM source WHERE name = 'manual'").Scan(&srcID); err != nil {
return nil, err
}
var im ProductImage
err := s.pool.QueryRow(ctx, `
INSERT INTO product_image (product_id, url, kind, license, source_id)
VALUES ($1,$2,$3,NULL,$4) RETURNING id, url, kind, license`,
productID, url, kind, srcID).Scan(&im.ID, &im.URL, &im.Kind, &im.License)
if err != nil {
return nil, err
}
if _, err := s.recomputeQuality(ctx, productID); err != nil {
return nil, err
}
_ = s.writeAudit(ctx, actor, "add_image", "product", &productID, []string{"image"}, nil, im)
return &im, nil
}
// DeleteImage removes an image and recomputes quality.
func (s *Store) DeleteImage(ctx context.Context, productID, imageID, actor string) error {
ct, err := s.pool.Exec(ctx, "DELETE FROM product_image WHERE id=$1 AND product_id=$2", imageID, productID)
if err != nil {
return err
}
if ct.RowsAffected() == 0 {
return ErrNotFound
}
if _, err := s.recomputeQuality(ctx, productID); err != nil {
return err
}
_ = s.writeAudit(ctx, actor, "delete_image", "product", &productID, []string{"image"}, map[string]string{"image_id": imageID}, nil)
return nil
}
// ---------- msrp ----------
// MSRPInput is the payload for adding an MSRP snapshot.
type MSRPInput 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"`
}
// AddMSRP inserts a suggested-retail-price snapshot.
func (s *Store) AddMSRP(ctx context.Context, productID, actor string, in MSRPInput) (*MSRP, error) {
if in.Currency == "" {
in.Currency = "CNY"
}
if in.Region == "" {
in.Region = "CN"
}
var srcID string
_ = s.pool.QueryRow(ctx, "SELECT id FROM source WHERE name = 'manual'").Scan(&srcID)
var m MSRP
err := s.pool.QueryRow(ctx, `
INSERT INTO product_msrp (product_id, amount, currency, region, source_id, source_url, effective_date, note)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8)
RETURNING id, amount, currency, region, effective_date::text, source_url, note`,
productID, in.Amount, in.Currency, in.Region, srcID, in.SourceURL, in.EffectiveDate, in.Note).
Scan(&m.ID, &m.Amount, &m.Currency, &m.Region, &m.EffectiveDate, &m.SourceURL, &m.Note)
if err != nil {
return nil, err
}
_ = s.writeAudit(ctx, actor, "add_msrp", "product", &productID, []string{"msrp"}, nil, m)
return &m, nil
}
// DeleteMSRP removes an MSRP snapshot.
func (s *Store) DeleteMSRP(ctx context.Context, productID, msrpID, actor string) error {
ct, err := s.pool.Exec(ctx, "DELETE FROM product_msrp WHERE id=$1 AND product_id=$2", msrpID, productID)
if err != nil {
return err
}
if ct.RowsAffected() == 0 {
return ErrNotFound
}
_ = s.writeAudit(ctx, actor, "delete_msrp", "product", &productID, []string{"msrp"}, map[string]string{"msrp_id": msrpID}, nil)
return nil
}
// ---------- dictionaries ----------
// Brand is a brand option for the edit form.
type Brand struct {
ID string `json:"id"`
Name string `json:"name"`
}
// ListBrands returns all brands ordered by name.
func (s *Store) ListBrands(ctx context.Context) ([]Brand, error) {
rows, err := s.pool.Query(ctx, "SELECT id, name FROM brand ORDER BY name")
if err != nil {
return nil, err
}
defer rows.Close()
out := []Brand{}
for rows.Next() {
var b Brand
if err := rows.Scan(&b.ID, &b.Name); err != nil {
return nil, err
}
out = append(out, b)
}
return out, rows.Err()
}
// Category is a category option for the edit form.
type Category struct {
ID string `json:"id"`
NameZH string `json:"name_zh"`
NameEN *string `json:"name_en"`
Path string `json:"path"`
Level int `json:"level"`
}
// ListCategories returns the full category tree.
func (s *Store) ListCategories(ctx context.Context) ([]Category, error) {
rows, err := s.pool.Query(ctx,
"SELECT id, name_zh, name_en, path::text, 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.Level); err != nil {
return nil, err
}
out = append(out, c)
}
return out, rows.Err()
}
// ---------- audit ----------
// AuditEntry is one audit-log row for the history view.
type AuditEntry struct {
ID string `json:"id"`
Actor string `json:"actor"`
Action string `json:"action"`
Fields []string `json:"fields"`
CreatedAt string `json:"created_at"`
}
// ListAudit returns audit history for one product, newest first.
func (s *Store) ListAudit(ctx context.Context, productID string, limit int) ([]AuditEntry, error) {
rows, err := s.pool.Query(ctx, `
SELECT id, actor, action, fields, created_at::text
FROM audit_log WHERE entity='product' AND entity_id=$1
ORDER BY created_at DESC LIMIT $2`, productID, limit)
if err != nil {
return nil, err
}
defer rows.Close()
out := []AuditEntry{}
for rows.Next() {
var e AuditEntry
if err := rows.Scan(&e.ID, &e.Actor, &e.Action, &e.Fields, &e.CreatedAt); err != nil {
return nil, err
}
out = append(out, e)
}
return out, rows.Err()
}
+7
View File
@@ -0,0 +1,7 @@
<!doctype html>
<html lang="zh">
<head><meta charset="utf-8" /><title>OpenGoods 管理后台</title></head>
<body>
<p>管理后台前端尚未构建。Docker 构建会在此处放入真正的前端产物。</p>
</body>
</html>
+21
View File
@@ -0,0 +1,21 @@
// Package adminweb embeds the built admin SPA (Vite dist). During Docker builds
// the real dist/ is produced by the node stage and copied in before go build;
// the committed placeholder keeps the package compilable for `go build ./...`.
package adminweb
import (
"embed"
"io/fs"
)
//go:embed all:dist
var distFS embed.FS
// Dist returns the embedded SPA filesystem rooted at dist/.
func Dist() fs.FS {
sub, err := fs.Sub(distFS, "dist")
if err != nil {
panic(err)
}
return sub
}
+126
View File
@@ -0,0 +1,126 @@
// Package auth provides minimal single-account authentication for the admin
// console: a bcrypt-verified login and a stdlib HMAC-SHA256 signed token
// (JWT-compatible) plus a chi middleware that guards write routes.
package auth
import (
"context"
"crypto/hmac"
"crypto/sha256"
"encoding/base64"
"encoding/json"
"errors"
"net/http"
"strings"
"time"
"golang.org/x/crypto/bcrypt"
)
// Authenticator holds the single admin credential and token signing secret.
type Authenticator struct {
username string
passwordHash []byte
secret []byte
ttl time.Duration
}
// New builds an Authenticator. passwordHash must be a bcrypt hash.
func New(username string, passwordHash, secret []byte, ttl time.Duration) *Authenticator {
return &Authenticator{username: username, passwordHash: passwordHash, secret: secret, ttl: ttl}
}
// ErrInvalidCredentials is returned when login fails.
var ErrInvalidCredentials = errors.New("invalid credentials")
// Login verifies the username/password and returns a signed token on success.
func (a *Authenticator) Login(username, password string) (string, error) {
if username != a.username {
// Still run bcrypt to keep timing roughly constant.
_ = bcrypt.CompareHashAndPassword(a.passwordHash, []byte(password))
return "", ErrInvalidCredentials
}
if err := bcrypt.CompareHashAndPassword(a.passwordHash, []byte(password)); err != nil {
return "", ErrInvalidCredentials
}
return a.issue(username)
}
type claims struct {
Sub string `json:"sub"`
Exp int64 `json:"exp"`
}
func b64(b []byte) string { return base64.RawURLEncoding.EncodeToString(b) }
func (a *Authenticator) sign(signingInput string) string {
mac := hmac.New(sha256.New, a.secret)
mac.Write([]byte(signingInput))
return b64(mac.Sum(nil))
}
func (a *Authenticator) issue(sub string) (string, error) {
header := b64([]byte(`{"alg":"HS256","typ":"JWT"}`))
payloadJSON, err := json.Marshal(claims{Sub: sub, Exp: time.Now().Add(a.ttl).Unix()})
if err != nil {
return "", err
}
payload := b64(payloadJSON)
signingInput := header + "." + payload
return signingInput + "." + a.sign(signingInput), nil
}
// Verify checks a token's signature and expiry, returning the subject.
func (a *Authenticator) Verify(token string) (string, error) {
parts := strings.Split(token, ".")
if len(parts) != 3 {
return "", errors.New("malformed token")
}
signingInput := parts[0] + "." + parts[1]
if !hmac.Equal([]byte(a.sign(signingInput)), []byte(parts[2])) {
return "", errors.New("bad signature")
}
payload, err := base64.RawURLEncoding.DecodeString(parts[1])
if err != nil {
return "", err
}
var c claims
if err := json.Unmarshal(payload, &c); err != nil {
return "", err
}
if time.Now().Unix() >= c.Exp {
return "", errors.New("token expired")
}
return c.Sub, nil
}
type ctxKey int
const userKey ctxKey = 0
// UserFrom returns the authenticated subject from the request context.
func UserFrom(ctx context.Context) string {
if v, ok := ctx.Value(userKey).(string); ok {
return v
}
return ""
}
// Middleware rejects requests without a valid Bearer token.
func (a *Authenticator) Middleware(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
h := r.Header.Get("Authorization")
token := strings.TrimPrefix(h, "Bearer ")
if token == h || token == "" {
http.Error(w, `{"error":{"code":"unauthorized","message":"missing token"}}`, http.StatusUnauthorized)
return
}
sub, err := a.Verify(token)
if err != nil {
http.Error(w, `{"error":{"code":"unauthorized","message":"invalid token"}}`, http.StatusUnauthorized)
return
}
ctx := context.WithValue(r.Context(), userKey, sub)
next.ServeHTTP(w, r.WithContext(ctx))
})
}
+62
View File
@@ -0,0 +1,62 @@
package auth
import (
"testing"
"time"
"golang.org/x/crypto/bcrypt"
)
func newTestAuth(t *testing.T, ttl time.Duration) *Authenticator {
t.Helper()
hash, err := bcrypt.GenerateFromPassword([]byte("s3cret"), bcrypt.MinCost)
if err != nil {
t.Fatalf("hash: %v", err)
}
return New("admin", hash, []byte("test-secret"), ttl)
}
func TestLoginAndVerify(t *testing.T) {
a := newTestAuth(t, time.Hour)
token, err := a.Login("admin", "s3cret")
if err != nil {
t.Fatalf("login: %v", err)
}
sub, err := a.Verify(token)
if err != nil {
t.Fatalf("verify: %v", err)
}
if sub != "admin" {
t.Fatalf("sub = %q, want admin", sub)
}
}
func TestLoginWrongCredentials(t *testing.T) {
a := newTestAuth(t, time.Hour)
if _, err := a.Login("admin", "nope"); err == nil {
t.Fatal("expected error for wrong password")
}
if _, err := a.Login("other", "s3cret"); err == nil {
t.Fatal("expected error for wrong username")
}
}
func TestVerifyRejectsTampered(t *testing.T) {
a := newTestAuth(t, time.Hour)
token, _ := a.Login("admin", "s3cret")
if _, err := a.Verify(token + "x"); err == nil {
t.Fatal("expected bad signature error")
}
if _, err := a.Verify("not.a.token"); err == nil {
t.Fatal("expected malformed/decoding error")
}
}
func TestVerifyRejectsExpired(t *testing.T) {
a := newTestAuth(t, -time.Minute)
token, _ := a.Login("admin", "s3cret")
if _, err := a.Verify(token); err == nil {
t.Fatal("expected expired token error")
}
}
+183 -15
View File
@@ -1,5 +1,4 @@
// Package handler wires up the public, read-only OpenGoods HTTP API.
//
// The OpenGoods service is a public-good product information API: it only
// collects and serves product facts. It exposes no purchase, checkout, or
// commerce endpoints by design.
@@ -7,48 +6,217 @@ package handler
import (
"encoding/json"
"errors"
"io/fs"
"net/http"
"strconv"
"strings"
"github.com/go-chi/chi/v5"
"github.com/go-chi/chi/v5/middleware"
"github.com/baicai2026-baicai/goods/api/internal/store"
)
// APIVersion is the current public API version prefix.
const APIVersion = "v1"
const (
defaultPageSize = 20
maxPageSize = 100
)
// Handler holds dependencies shared by the HTTP routes.
type Handler struct {
store *store.Store
spa fs.FS
}
// New constructs a Handler backed by the given store. spa may be nil (JSON-only).
func New(s *store.Store, spa fs.FS) *Handler {
return &Handler{store: s, spa: spa}
}
// Router builds the top-level HTTP handler with middleware and routes mounted.
func Router() http.Handler {
func (h *Handler) Router() http.Handler {
r := chi.NewRouter()
r.Use(middleware.RequestID)
r.Use(middleware.RealIP)
r.Use(middleware.Recoverer)
r.Get("/healthz", Healthz)
r.Get("/healthz", h.Healthz)
r.Route("/api/"+APIVersion, func(r chi.Router) {
r.Route("/products", func(r chi.Router) {
r.Get("/barcode/{gtin}", notImplemented)
r.Get("/search", notImplemented)
r.Get("/{id}", notImplemented)
r.Get("/{id}/nutriments", notImplemented)
r.Get("/{id}/msrp", notImplemented)
r.Get("/barcode/{gtin}", h.ProductByBarcode)
r.Get("/search", h.SearchProducts)
r.Get("/{id}", h.ProductByID)
r.Get("/{id}/nutriments", h.ProductNutriments)
r.Get("/{id}/msrp", h.ProductMSRP)
})
r.Get("/brands", notImplemented)
r.Get("/categories", notImplemented)
r.Get("/sources/{id}", notImplemented)
r.Get("/brands", h.ListBrands)
r.Get("/categories", h.ListCategories)
r.Get("/sources/{id}", h.SourceByID)
})
// Public SPA (homepage + search + contribute). API routes above take
// precedence; everything else falls back to the embedded single-page app.
if h.spa != nil {
r.Handle("/*", http.HandlerFunc(h.serveSPA))
}
return r
}
func (h *Handler) serveSPA(w http.ResponseWriter, r *http.Request) {
rel := strings.TrimPrefix(r.URL.Path, "/")
if rel == "" {
rel = "index.html"
}
if f, err := h.spa.Open(rel); err == nil {
f.Close()
http.FileServer(http.FS(h.spa)).ServeHTTP(w, r)
return
}
// SPA fallback: serve index.html for client-side routes.
data, err := fs.ReadFile(h.spa, "index.html")
if err != nil {
http.NotFound(w, r)
return
}
w.Header().Set("Content-Type", "text/html; charset=utf-8")
_, _ = w.Write(data)
}
// Healthz reports liveness of the service.
func Healthz(w http.ResponseWriter, r *http.Request) {
func (h *Handler) Healthz(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, map[string]string{"status": "ok"})
}
// notImplemented is a placeholder for endpoints scoped to later milestones.
func notImplemented(w http.ResponseWriter, r *http.Request) {
writeError(w, r, http.StatusNotImplemented, "not_implemented", "endpoint not implemented yet")
// ProductByBarcode returns a product by its GTIN.
func (h *Handler) ProductByBarcode(w http.ResponseWriter, r *http.Request) {
p, err := h.store.ProductByGTIN(r.Context(), chi.URLParam(r, "gtin"))
if h.handleErr(w, r, err) {
return
}
writeJSON(w, http.StatusOK, p)
}
// ProductByID returns a product by its UUID.
func (h *Handler) ProductByID(w http.ResponseWriter, r *http.Request) {
p, err := h.store.ProductByID(r.Context(), chi.URLParam(r, "id"))
if h.handleErr(w, r, err) {
return
}
writeJSON(w, http.StatusOK, p)
}
// SearchProducts runs a fuzzy name search with optional category filter + paging.
func (h *Handler) SearchProducts(w http.ResponseWriter, r *http.Request) {
q := r.URL.Query().Get("q")
category := r.URL.Query().Get("category")
page, size := pageParams(r)
items, total, err := h.store.SearchProducts(r.Context(), q, category, size, (page-1)*size)
if h.handleErr(w, r, err) {
return
}
writeJSON(w, http.StatusOK, map[string]any{
"items": items,
"page": page,
"size": size,
"total": total,
})
}
// ProductNutriments returns just the nutrition facts of a product.
func (h *Handler) ProductNutriments(w http.ResponseWriter, r *http.Request) {
n, err := h.store.Nutriments(r.Context(), chi.URLParam(r, "id"))
if h.handleErr(w, r, err) {
return
}
writeJSON(w, http.StatusOK, n)
}
// ProductMSRP returns official suggested retail price snapshots (no purchase link).
func (h *Handler) ProductMSRP(w http.ResponseWriter, r *http.Request) {
items, err := h.store.ListMSRP(r.Context(), chi.URLParam(r, "id"))
if h.handleErr(w, r, err) {
return
}
writeJSON(w, http.StatusOK, map[string]any{
"items": items,
"disclaimer": "厂商建议零售价历史快照,仅供参考,不构成购买建议,本服务不提供任何购买入口。",
})
}
// ListBrands returns a paginated list of brands.
func (h *Handler) ListBrands(w http.ResponseWriter, r *http.Request) {
page, size := pageParams(r)
items, total, err := h.store.ListBrands(r.Context(), size, (page-1)*size)
if h.handleErr(w, r, err) {
return
}
writeJSON(w, http.StatusOK, map[string]any{
"items": items, "page": page, "size": size, "total": total,
})
}
// ListCategories returns the full category tree.
func (h *Handler) ListCategories(w http.ResponseWriter, r *http.Request) {
items, err := h.store.ListCategories(r.Context())
if h.handleErr(w, r, err) {
return
}
writeJSON(w, http.StatusOK, map[string]any{"items": items})
}
// SourceByID returns a single data source.
func (h *Handler) SourceByID(w http.ResponseWriter, r *http.Request) {
src, err := h.store.SourceByID(r.Context(), chi.URLParam(r, "id"))
if h.handleErr(w, r, err) {
return
}
writeJSON(w, http.StatusOK, src)
}
// handleErr writes an appropriate error response; returns true if it handled one.
func (h *Handler) handleErr(w http.ResponseWriter, r *http.Request, err error) bool {
if err == nil {
return false
}
if errors.Is(err, store.ErrNotFound) {
writeError(w, r, http.StatusNotFound, "not_found", "resource not found")
return true
}
writeError(w, r, http.StatusInternalServerError, "internal_error", "internal server error")
return true
}
func pageParams(r *http.Request) (page, size int) {
page = atoiDefault(r.URL.Query().Get("page"), 1)
if page < 1 {
page = 1
}
size = atoiDefault(r.URL.Query().Get("size"), defaultPageSize)
if size < 1 {
size = defaultPageSize
}
if size > maxPageSize {
size = maxPageSize
}
return page, size
}
func atoiDefault(s string, fallback int) int {
if s == "" {
return fallback
}
v, err := strconv.Atoi(s)
if err != nil {
return fallback
}
return v
}
func writeJSON(w http.ResponseWriter, status int, body any) {
+134
View File
@@ -0,0 +1,134 @@
package handler
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/baicai2026-baicai/goods/api/internal/store"
)
// newTestHandler connects to the test database, skipping if unavailable or
// unmigrated. It inserts a known product (cleaned up via t.Cleanup) so the
// endpoint assertions are deterministic.
func newTestHandler(t *testing.T) (*Handler, string) {
t.Helper()
dsn := os.Getenv("OPENGOODS_DATABASE_URL")
if dsn == "" {
dsn = "postgres://opengoods:opengoods@localhost:5432/opengoods?sslmode=disable"
}
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Skipf("no database: %v", err)
}
if err := pool.Ping(ctx); err != nil {
pool.Close()
t.Skipf("database not reachable: %v", err)
}
var hasProduct bool
if err := pool.QueryRow(ctx, "SELECT to_regclass('public.product') IS NOT NULL").Scan(&hasProduct); err != nil || !hasProduct {
pool.Close()
t.Skip("migrations not applied")
}
gtin := "4006381333931"
_, err = pool.Exec(context.Background(), `
INSERT INTO product (gtin, name, category_id, net_content_value, net_content_unit)
VALUES ($1, 'Test Cola', (SELECT id FROM category WHERE path='food.beverages.carbonated'), 330, 'ml')
ON CONFLICT (gtin) WHERE gtin IS NOT NULL DO UPDATE SET name = EXCLUDED.name`, gtin)
if err != nil {
pool.Close()
t.Fatalf("seed insert failed: %v", err)
}
var pid string
_ = pool.QueryRow(context.Background(), "SELECT id FROM product WHERE gtin=$1", gtin).Scan(&pid)
_, _ = pool.Exec(context.Background(), `
INSERT INTO food_detail (product_id, nutrition_basis, nutriments)
VALUES ($1, 'per_100ml', '{"energy_kcal": 42}'::jsonb)
ON CONFLICT (product_id) DO UPDATE SET nutriments = EXCLUDED.nutriments`, pid)
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), "DELETE FROM product WHERE gtin=$1", gtin)
pool.Close()
})
return New(store.New(pool), nil), gtin
}
func doGET(t *testing.T, h *Handler, path string) *httptest.ResponseRecorder {
t.Helper()
req := httptest.NewRequest(http.MethodGet, path, nil)
rec := httptest.NewRecorder()
h.Router().ServeHTTP(rec, req)
return rec
}
func TestProductByBarcode(t *testing.T) {
h, gtin := newTestHandler(t)
rec := doGET(t, h, "/api/"+APIVersion+"/products/barcode/"+gtin)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
var p store.Product
if err := json.NewDecoder(rec.Body).Decode(&p); err != nil {
t.Fatal(err)
}
if p.Name != "Test Cola" || p.GTIN == nil || *p.GTIN != gtin {
t.Fatalf("unexpected product: %+v", p)
}
if p.CategoryPath == nil || *p.CategoryPath != "food.beverages.carbonated" {
t.Fatalf("category not joined: %+v", p.CategoryPath)
}
}
func TestProductByBarcodeNotFound(t *testing.T) {
h, _ := newTestHandler(t)
rec := doGET(t, h, "/api/"+APIVersion+"/products/barcode/0000000000000")
if rec.Code != http.StatusNotFound {
t.Fatalf("expected 404, got %d", rec.Code)
}
}
func TestSearchProducts(t *testing.T) {
h, _ := newTestHandler(t)
rec := doGET(t, h, "/api/"+APIVersion+"/products/search?q=Cola&category=food.beverages")
if rec.Code != http.StatusOK {
t.Fatalf("status = %d", rec.Code)
}
var body struct {
Items []store.ProductSummary `json:"items"`
Total int `json:"total"`
}
if err := json.NewDecoder(rec.Body).Decode(&body); err != nil {
t.Fatal(err)
}
if body.Total < 1 {
t.Fatalf("expected at least 1 result, got %d", body.Total)
}
}
func TestListCategories(t *testing.T) {
h, _ := newTestHandler(t)
rec := doGET(t, h, "/api/"+APIVersion+"/categories")
if rec.Code != http.StatusOK {
t.Fatalf("status = %d", rec.Code)
}
var body struct {
Items []store.Category `json:"items"`
}
if err := json.NewDecoder(rec.Body).Decode(&body); err != nil {
t.Fatal(err)
}
if len(body.Items) < 20 {
t.Fatalf("expected seeded categories, got %d", len(body.Items))
}
}
+19 -9
View File
@@ -11,7 +11,7 @@ func TestHealthz(t *testing.T) {
req := httptest.NewRequest(http.MethodGet, "/healthz", nil)
rec := httptest.NewRecorder()
Router().ServeHTTP(rec, req)
New(nil, nil).Router().ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("expected status %d, got %d", http.StatusOK, rec.Code)
@@ -26,13 +26,23 @@ func TestHealthz(t *testing.T) {
}
}
func TestProductEndpointNotImplemented(t *testing.T) {
req := httptest.NewRequest(http.MethodGet, "/api/"+APIVersion+"/products/barcode/3017624010701", nil)
rec := httptest.NewRecorder()
Router().ServeHTTP(rec, req)
if rec.Code != http.StatusNotImplemented {
t.Fatalf("expected status %d, got %d", http.StatusNotImplemented, rec.Code)
func TestPageParams(t *testing.T) {
cases := []struct {
query string
wantPage, wantSz int
}{
{"", 1, defaultPageSize},
{"page=3&size=10", 3, 10},
{"page=0&size=-5", 1, defaultPageSize},
{"size=1000", 1, maxPageSize},
{"page=abc", 1, defaultPageSize},
}
for _, c := range cases {
req := httptest.NewRequest(http.MethodGet, "/?"+c.query, nil)
page, size := pageParams(req)
if page != c.wantPage || size != c.wantSz {
t.Errorf("query %q: got page=%d size=%d, want page=%d size=%d",
c.query, page, size, c.wantPage, c.wantSz)
}
}
}
+10
View File
@@ -0,0 +1,10 @@
<!doctype html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8" />
<title>OpenGoods</title>
</head>
<body>
<div id="root">OpenGoods public site placeholder. Built assets are injected during Docker build.</div>
</body>
</html>
+22
View File
@@ -0,0 +1,22 @@
// Package publicweb embeds the built public SPA (Vite dist). During Docker
// builds the real dist/ is produced by the node stage and copied in before go
// build; the committed placeholder keeps the package compilable for
// `go build ./...`.
package publicweb
import (
"embed"
"io/fs"
)
//go:embed all:dist
var distFS embed.FS
// Dist returns the embedded SPA filesystem rooted at dist/.
func Dist() fs.FS {
sub, err := fs.Sub(distFS, "dist")
if err != nil {
panic(err)
}
return sub
}
+278
View File
@@ -0,0 +1,278 @@
// 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"
"errors"
"strconv"
"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)
}
// 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"`
NetContentValue *float64 `json:"net_content_value"`
NetContentUnit *string `json:"net_content_unit"`
CountryOfOrigin *string `json:"country_of_origin"`
QualityScore float64 `json:"quality_score"`
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"`
}
// 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"`
}
const productSelect = `
SELECT p.id, p.gtin, p.name, b.name, c.path::text, p.gpc_brick_code,
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, error) {
var p Product
err := row.Scan(
&p.ID, &p.GTIN, &p.Name, &p.Brand, &p.CategoryPath, &p.GPCBrickCode,
&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, ErrNotFound
}
if err != nil {
return nil, err
}
return &p, nil
}
// ProductByGTIN looks up an active product by its barcode.
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)
}
// 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)
}
// SearchProducts performs a fuzzy name search with optional category subtree filter.
func (s *Store) SearchProducts(ctx context.Context, q, category string, limit, offset int) ([]ProductSummary, int, error) {
args := []any{}
where := "WHERE p.status = 'active'"
if q != "" {
args = append(args, q)
where += " AND p.name ILIKE '%' || $1 || '%'"
}
if category != "" {
args = append(args, category)
where += " AND c.path <@ $" + strconv.Itoa(len(args)) + "::ltree"
}
countSQL := "SELECT count(*) FROM product p LEFT JOIN category c ON c.id = p.category_id " + where
var total int
if err := s.pool.QueryRow(ctx, countSQL, args...).Scan(&total); err != nil {
return nil, 0, err
}
args = append(args, limit, offset)
listSQL := `
SELECT p.id, p.gtin, p.name, b.name, c.path::text
FROM product p
LEFT JOIN brand b ON b.id = p.brand_id
LEFT JOIN category c ON c.id = p.category_id ` + where +
" ORDER BY p.name 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); 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()
}
// 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
}
+79
View File
@@ -0,0 +1,79 @@
name: goods
services:
postgres:
image: postgres:16-alpine
restart: unless-stopped
environment:
POSTGRES_USER: ${POSTGRES_USER}
POSTGRES_PASSWORD: ${POSTGRES_PASSWORD}
POSTGRES_DB: ${POSTGRES_DB}
volumes:
- pgdata:/var/lib/postgresql/data
healthcheck:
test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER}"]
interval: 5s
timeout: 5s
retries: 10
redis:
image: redis:7-alpine
restart: unless-stopped
healthcheck:
test: ["CMD", "redis-cli", "ping"]
interval: 5s
timeout: 5s
retries: 10
minio:
image: minio/minio:latest
restart: unless-stopped
command: server /data --console-address ":9001"
environment:
MINIO_ROOT_USER: ${MINIO_ROOT_USER}
MINIO_ROOT_PASSWORD: ${MINIO_ROOT_PASSWORD}
volumes:
- miniodata:/data
healthcheck:
test: ["CMD", "mc", "ready", "local"]
interval: 5s
timeout: 5s
retries: 10
api:
build:
context: .
dockerfile: api/Dockerfile.prod
restart: unless-stopped
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
environment:
OPENGOODS_ADDR: ":8080"
OPENGOODS_DATABASE_URL: "postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB}?sslmode=disable"
OPENGOODS_REDIS_URL: "redis://redis:6379/0"
ports:
- "127.0.0.1:8120:8080"
admin:
build:
context: .
dockerfile: api/Dockerfile.admin
restart: unless-stopped
depends_on:
postgres:
condition: service_healthy
environment:
GOODS_ADMIN_ADDR: ":8080"
OPENGOODS_DATABASE_URL: "postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB}?sslmode=disable"
GOODS_ADMIN_BASE_PATH: "/ping"
GOODS_ADMIN_USER: "${GOODS_ADMIN_USER}"
GOODS_ADMIN_PASSWORD: "${GOODS_ADMIN_PASSWORD}"
GOODS_ADMIN_JWT_SECRET: "${GOODS_ADMIN_JWT_SECRET}"
ports:
- "127.0.0.1:8121:8080"
volumes:
pgdata:
miniodata:
+34
View File
@@ -0,0 +1,34 @@
# 生产部署 (Docker)
`docker-compose.prod.yml` 部署,与本地 `docker-compose.yml` 的区别:
-`api` 映射宿主端口,且绑定 `127.0.0.1:8120`(由外层 nginx 反代 + HTTPS);`postgres`/`redis`/`minio` 不对外暴露端口,仅容器内网互通。
- 所有服务 `restart: unless-stopped`
- 凭据从 `.env` 注入(见 `.env.example`),不写入仓库。
- `api` 使用 `api/Dockerfile.prod`:运行镜像用 `scratch`(从构建镜像拷贝 ca-certs),适用于 `gcr.io/distroless` 不可达的环境;Go 模块走 `goproxy.cn`
- `admin`(运营后台):带登录的写入服务 + 内嵌前端,绑定 `127.0.0.1:8121`,由 nginx 反代到公开站点的 `/ping` 路径。镜像 `api/Dockerfile.admin`(node 构建前端 → 内嵌进 Go 二进制 → scratch 运行)。仅 `admin` 可写库(人工编辑以 `source=manual` 记录字段级溯源 + `audit_log` 留痕),公开 `api` 仍只读。
## 步骤
```bash
cp .env.example .env # 填入真实随机密码
docker compose -f docker-compose.prod.yml up -d --build
# 迁移(migrate 容器接入同一网络,DSN 指向 postgres 服务)
set -a; . ./.env; set +a
DBURL="postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB}?sslmode=disable"
docker run --rm --network goods_default -v "$PWD/migrations:/migrations" \
migrate/migrate -path=/migrations -database "$DBURL" up
curl -s http://127.0.0.1:8120/healthz # {"status":"ok"}
```
nginx 反代(子域 + HTTPS):80 端口 301 跳转到 443443 `proxy_pass http://127.0.0.1:8120`,证书用 acme.sh 签发并配 `--reloadcmd "nginx -s reload"` 自动续期。
运营后台 `/ping`(同域复用证书):在 443 server 块内加一段
```nginx
location /ping { proxy_pass http://127.0.0.1:8121; }
```
后台凭据见 `.env``GOODS_ADMIN_USER` / `GOODS_ADMIN_PASSWORD` / `GOODS_ADMIN_JWT_SECRET`。新增迁移 `0005_admin``audit_log` 表 + `manual` 来源)随 `migrate ... up` 自动应用。
+76
View File
@@ -0,0 +1,76 @@
# 采集管理 (M4)
M4 在 M2Open Food Facts 首次导入)基础上,补齐"持续运营"所需的采集能力:
增量更新、第二数据源补全(GS1)、去重合并与字段级冲突解决、数据质量评分,
以及把这些串起来的定时调度。全部为 Python 侧(`ingestion/`),只写库、可单测。
## 组成
| 能力 | 模块 | 说明 |
|------|------|------|
| 增量采集 | `adapters/openfoodfacts.py: fetch_modified_since()` | 按 `last_modified_t` 拉取自上次水位后变更的商品 |
| 采集水位 | `etl/state.py` + `ingest_state` 表 | 每个源持久化 `last_modified_t`,只前进不回退 |
| GS1 补全 | `adapters/gs1.py` + `etl/supplement.py` | 用权威条码源补**缺失**字段(品牌/厂商/GPC/产地/净含量),不覆盖已有值 |
| 去重合并 | `etl/dedup.py` | 非 GTIN 重复(同名+品牌+净含量)合并到质量最高的主记录 |
| 冲突解决 | `etl/merge.py` | 多源同字段按"源权重 > 新鲜度"择优,保留字段级溯源 |
| 质量评分 | `etl/quality.py` | 0~1 分,落到 `product.quality_score` |
| 定时调度 | `jobs/schedule.py` | 固定周期跑"增量 + 去重"一轮,零额外依赖 |
## 质量评分
锁定公式(各分量均归一到 0~1):
```
quality = 0.4 * 完整度 + 0.3 * 源权重 + 0.2 * 多源一致 + 0.1 * 新鲜度
```
- **完整度**`name/gtin/brand/category/net_content/country/nutriments/ingredients/image` 9 项的命中比例。
- **源权重**:贡献该商品的源中最高 `source.trust_weight`OFF=0.7GS1=0.9)。
- **多源一致**:源数量代理——单源 0.5、两源 0.8、三源及以上 1.0(单源无法互证)。
- **新鲜度**:最近一次 `product_source.fetched_at` 的时间衰减(≤30d=1.0 … >730d=0.2)。
`load_record()``merge_products()` 写入后都会调 `update_quality()` 重算。
## 增量水位
`ingest_state`(迁移 `0004`)每源一行,记录 `last_modified_t``last_run_at``stats`
`set_watermark()``GREATEST(...)` 保证水位只前进,避免乱序/中断的运行回退进度。
## 运行
前置:`docker compose up -d postgres` 且迁移已 `up`(含 `0004`)。DSN 默认读 `OPENGOODS_DATABASE_URL`
```bash
# 增量更新 OFF(从持久化水位开始;--since 可覆盖)
python -m opengoods.jobs.update_off --max-pages 5
python -m opengoods.jobs.update_off --since 1700000000
# 去重合并(--dry-run 只报告不写库)
python -m opengoods.jobs.dedup --dry-run
python -m opengoods.jobs.dedup --actor nightly
# 定时调度:单轮 / 周期循环(增量 + 去重)
python -m opengoods.jobs.schedule --once
python -m opengoods.jobs.schedule --interval 3600
```
## GS1 补全
GS1 为付费、分区域的授权数据,适配器支持两种模式:
- **离线**(默认):从本地 JSON 映射 `{gtin: {...}}` 查(`GS1Adapter.from_file(path)`),
供测试与内网环境使用。
- **在线**:传 `base_url` + `client`+ `api_key`),`GET {base_url}/{gtin}`,按
Verified-by-GS1 风格字段解析。
补全只填**空缺**字段并在 `product_source` 记字段级溯源,源标记为 `gs1`
## 测试
```bash
cd ingestion && pip install -e ".[dev]"
ruff check . && ruff format --check . && pytest -q
```
纯函数测试(质量/冲突/增量分页)始终运行;依赖库的测试(水位/质量落库/去重/GS1 补全)
在无数据库或未应用 M4 迁移时自动跳过。
+124
View File
@@ -0,0 +1,124 @@
"""GS1 barcode supplement adapter.
GS1 (e.g. *Verified by GS1* / GS1 China) is the authoritative registry that maps
a GTIN to its brand owner, product description and GPC category. We use it to
*supplement* — fill gaps in — records gathered from crowd sources like Open Food
Facts, never to overwrite existing values.
Real GS1 access is credentialed and region-specific, so this adapter supports
two modes:
* **offline** (default): look barcodes up in a local JSON mapping file. This is
what tests and air-gapped runs use.
* **online**: GET ``{base_url}/{gtin}`` with an API key header, then normalize
the response. Enabled by passing ``base_url`` + ``client``.
Either way :meth:`fetch_barcode` returns a normalized *supplement* dict (or
``None``); :mod:`opengoods.etl.supplement` applies it to the database.
"""
from __future__ import annotations
import json
from collections.abc import Iterator
from pathlib import Path
import httpx
SOURCE_NAME = "gs1"
GS1_HOMEPAGE = "https://www.gs1.org"
GS1_LICENSE = "proprietary"
# GS1 is the authoritative barcode registry -> high trust.
GS1_TRUST = 0.9
# Keys of a normalized supplement record.
_SUPPLEMENT_KEYS = (
"gtin",
"name",
"brand",
"manufacturer",
"gpc_brick_code",
"country_of_origin",
"net_content_value",
"net_content_unit",
)
def _normalize(code: str, data: dict) -> dict:
"""Project a raw mapping/record onto the supplement schema (non-empty only)."""
rec: dict = {"gtin": code}
for key in _SUPPLEMENT_KEYS:
if key == "gtin":
continue
value = data.get(key)
if value not in (None, "", []):
rec[key] = value
return rec
def _parse_api(code: str, payload: dict) -> dict:
"""Best-effort mapping of a Verified-by-GS1 style payload to our schema."""
item = payload
if isinstance(payload.get("gtinRecords"), list) and payload["gtinRecords"]:
item = payload["gtinRecords"][0]
return _normalize(
code,
{
"name": item.get("productDescription") or item.get("description"),
"brand": item.get("brandName"),
"manufacturer": item.get("companyName") or item.get("licenseeName"),
"gpc_brick_code": item.get("gpcCategoryCode"),
"country_of_origin": item.get("countryOfSaleCode") or item.get("countryCode"),
"net_content_value": item.get("netContent"),
"net_content_unit": item.get("netContentUnit"),
},
)
class GS1Adapter:
"""Look up GTIN supplements from a local mapping or a GS1-style API."""
source_name = SOURCE_NAME
def __init__(
self,
mapping: dict | None = None,
*,
client: httpx.Client | None = None,
base_url: str | None = None,
api_key: str | None = None,
) -> None:
self._mapping = mapping or {}
self._client = client
self._base_url = base_url.rstrip("/") if base_url else None
self._api_key = api_key
@classmethod
def from_file(cls, path: str | Path) -> GS1Adapter:
"""Build an offline adapter from a JSON ``{gtin: {...}}`` mapping file."""
data = json.loads(Path(path).read_text(encoding="utf-8"))
return cls(mapping=data)
def fetch_barcode(self, code: str) -> dict | None:
"""Return a normalized supplement dict for ``code`` (or ``None``)."""
if self._base_url and self._client is not None:
headers = {"apikey": self._api_key} if self._api_key else {}
resp = self._client.get(f"{self._base_url}/{code}", headers=headers)
if resp.status_code == 404:
return None
resp.raise_for_status()
rec = _parse_api(code, resp.json())
else:
data = self._mapping.get(code)
if not data:
return None
rec = _normalize(code, data)
# A record with only the GTIN carries no supplement.
return rec if len(rec) > 1 else None
def fetch(self, barcodes: list[str]) -> Iterator[dict]:
"""Yield supplement records for the given barcodes."""
for code in barcodes:
rec = self.fetch_barcode(code)
if rec is not None:
yield rec
@@ -25,6 +25,16 @@ USER_AGENT = "OpenGoods/0.1 (+https://github.com/baicai2026-baicai/goods) public
# Conservative client-side spacing between API calls (seconds).
_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"
# Fields requested from the search API so a returned product can be transformed
# without an extra per-barcode round trip.
_SEARCH_FIELDS = (
"code,product_name,product_name_en,product_name_zh,brands,quantity,"
"categories,categories_tags,countries,ingredients_text,allergens_tags,"
"additives_tags,nutriments,nutriscore_grade,serving_size,"
"image_front_url,image_url,last_modified_t"
)
class OpenFoodFactsAdapter:
@@ -65,6 +75,45 @@ class OpenFoodFactsAdapter:
if record is not None:
yield record
def fetch_modified_since(
self,
since_t: int,
*,
page_size: int = 100,
max_pages: int = 10,
) -> Iterator[dict]:
"""Yield products modified after ``since_t`` (unix ``last_modified_t``).
Uses the OFF search API sorted by ``last_modified_t`` (most recent
first) and paginates until it reaches products at or before the
watermark, an empty/short page, or ``max_pages``. This is the
incremental ingestion path: callers persist the highest
``last_modified_t`` they processed as the next watermark.
"""
for page in range(1, max_pages + 1):
self._throttle()
resp = self._client.get(
_SEARCH_URL,
params={
"fields": _SEARCH_FIELDS,
"sort_by": "last_modified_t",
"page": page,
"page_size": page_size,
},
)
resp.raise_for_status()
products = resp.json().get("products") or []
if not products:
return
reached_old = False
for prod in products:
if int(prod.get("last_modified_t") or 0) <= since_t:
reached_old = True
break
yield prod
if reached_old or len(products) < page_size:
return
def read_dump(path: str | Path) -> Iterator[dict]:
"""Yield raw product records from an OFF JSONL dump file.
+139
View File
@@ -0,0 +1,139 @@
"""Duplicate detection and product merging.
Barcodes (GTIN) are already unique at the schema level, so duplicates here are
non-GTIN records that describe the same product (same normalized name + brand +
net content). For each duplicate group we keep the highest-quality product as
canonical and merge the rest into it: child rows (provenance, images, MSRP) are
re-pointed to the canonical product, the merged product is marked ``merged``
with ``canonical_id`` set, and a row is written to ``merge_log``.
"""
from __future__ import annotations
from typing import Any
import psycopg
from opengoods.etl.quality import update_quality
def _norm(text: str | None) -> str:
return " ".join((text or "").lower().split())
def product_signature(name: str | None, brand: str | None, net_canonical: Any | None) -> str | None:
"""Stable signature for non-GTIN dedup, or ``None`` if too sparse to match."""
n = _norm(name)
if not n:
return None
net = "" if net_canonical is None else str(net_canonical)
return f"{n}|{_norm(brand)}|{net}"
def choose_canonical(members: list[dict]) -> dict:
"""Pick the canonical product: best quality, then oldest, then lowest id."""
return min(
members,
key=lambda m: (
-float(m.get("quality_score") or 0.0),
m.get("created_at"),
str(m.get("id")),
),
)
def find_duplicate_groups(conn: psycopg.Connection) -> list[list[dict]]:
"""Return groups (size >= 2) of active products sharing a signature."""
rows = conn.execute(
"""
SELECT p.id, p.name, b.normalized_name, p.net_content_canonical,
p.quality_score, p.created_at
FROM product p
LEFT JOIN brand b ON b.id = p.brand_id
WHERE p.status = 'active'
"""
).fetchall()
groups: dict[str, list[dict]] = {}
for r in rows:
sig = product_signature(r[1], r[2], r[3])
if sig is None:
continue
member = {
"id": r[0],
"name": r[1],
"quality_score": r[4],
"created_at": r[5],
}
groups.setdefault(sig, []).append(member)
return [m for m in groups.values() if len(m) >= 2]
def merge_products(
conn: psycopg.Connection,
kept_id: str,
merged_id: str,
reason: str = "auto-dedup",
actor: str = "ingestion",
) -> None:
"""Merge ``merged_id`` into ``kept_id`` (re-point children, mark merged)."""
if kept_id == merged_id:
return
# Re-point provenance, images and MSRP to the canonical product.
conn.execute(
"UPDATE product_source SET product_id = %s WHERE product_id = %s",
(kept_id, merged_id),
)
conn.execute(
"UPDATE product_image SET product_id = %s WHERE product_id = %s",
(kept_id, merged_id),
)
conn.execute(
"UPDATE product_msrp SET product_id = %s WHERE product_id = %s",
(kept_id, merged_id),
)
# food_detail has product_id as PK, so it can only move if the canonical
# product does not already have one.
kept_has_food = conn.execute(
"SELECT 1 FROM food_detail WHERE product_id = %s", (kept_id,)
).fetchone()
if not kept_has_food:
conn.execute(
"UPDATE food_detail SET product_id = %s WHERE product_id = %s",
(kept_id, merged_id),
)
conn.execute(
"UPDATE product SET status = 'merged', canonical_id = %s WHERE id = %s",
(kept_id, merged_id),
)
conn.execute(
"""
INSERT INTO merge_log (kept_id, merged_id, reason, actor)
VALUES (%s, %s, %s, %s)
""",
(kept_id, merged_id, reason, actor),
)
# The canonical product gained sources, so its quality may have changed.
update_quality(conn, kept_id)
def dedup_all(
conn: psycopg.Connection, actor: str = "ingestion", dry_run: bool = False
) -> dict[str, int]:
"""Merge every duplicate group. Returns counts of groups and merges."""
groups = find_duplicate_groups(conn)
merged = 0
for members in groups:
canonical = choose_canonical(members)
for m in members:
if m["id"] == canonical["id"]:
continue
if not dry_run:
merge_products(conn, canonical["id"], m["id"], actor=actor)
merged += 1
return {"groups": len(groups), "merged": merged}
+18 -3
View File
@@ -14,6 +14,7 @@ import psycopg
from psycopg.types.json import Jsonb
from opengoods.adapters.openfoodfacts import OFF_LICENSE, SOURCE_NAME
from opengoods.etl.quality import update_quality
OFF_HOMEPAGE = "https://world.openfoodfacts.org"
@@ -29,8 +30,14 @@ def _normalize_brand(name: str) -> str:
return " ".join(name.lower().split())
def ensure_source(conn: psycopg.Connection) -> str:
"""Upsert the Open Food Facts source row and return its id."""
def ensure_source_named(
conn: psycopg.Connection,
name: str,
homepage: str,
license: str,
trust_weight: float,
) -> str:
"""Upsert a source row by name and return its id."""
row = conn.execute(
"""
INSERT INTO source (name, homepage, license, trust_weight)
@@ -38,11 +45,16 @@ def ensure_source(conn: psycopg.Connection) -> str:
ON CONFLICT (name) DO UPDATE SET homepage = EXCLUDED.homepage
RETURNING id
""",
(SOURCE_NAME, OFF_HOMEPAGE, OFF_LICENSE, 0.7),
(name, homepage, license, trust_weight),
).fetchone()
return row[0]
def ensure_source(conn: psycopg.Connection) -> str:
"""Upsert the Open Food Facts source row and return its id."""
return ensure_source_named(conn, SOURCE_NAME, OFF_HOMEPAGE, OFF_LICENSE, 0.7)
def _ensure_brand(conn: psycopg.Connection, name: str | None) -> str | None:
if not name:
return None
@@ -178,6 +190,9 @@ def load_record(conn: psycopg.Connection, rec: dict[str, Any], source_id: str, r
Jsonb(_jsonable(raw)),
),
)
# Recompute the data-quality score now that all facts + provenance exist.
update_quality(conn, product_id)
return product_id
+104
View File
@@ -0,0 +1,104 @@
"""Field-level conflict resolution for multi-source records.
When more than one source describes the same product, each field may have
several candidate values. We pick a winner per field by source trust first,
then recency, ignoring empty values, and keep a provenance trail of which
source won each field.
These are pure functions (no DB / no network) so they are easy to unit-test;
the DB-level record merge lives in :mod:`opengoods.etl.dedup`.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from datetime import datetime
@dataclass(frozen=True)
class Candidate:
"""One source's proposed value for a field."""
value: object
source: str
trust: float = 0.5
fetched_at: datetime | None = None
@dataclass
class FieldResolution:
"""The winning value for a field plus the source it came from."""
value: object
source: str | None = None
@dataclass
class MergedRecord:
"""A merged record with per-field provenance (field name -> source)."""
values: dict[str, object] = field(default_factory=dict)
provenance: dict[str, str] = field(default_factory=dict)
def _is_empty(value: object) -> bool:
if value is None:
return True
if isinstance(value, str):
return value.strip() == ""
if isinstance(value, (list, dict, tuple, set)):
return len(value) == 0
return False
def _sort_key(c: Candidate) -> tuple[float, float]:
ts = c.fetched_at.timestamp() if c.fetched_at is not None else float("-inf")
return (c.trust, ts)
def resolve_field(candidates: list[Candidate]) -> FieldResolution | None:
"""Pick the best non-empty candidate for one field.
Ranking: highest source trust, then most recent ``fetched_at``. Returns
``None`` when there is no usable (non-empty) candidate.
"""
usable = [c for c in candidates if not _is_empty(c.value)]
if not usable:
return None
winner = max(usable, key=_sort_key)
return FieldResolution(value=winner.value, source=winner.source)
def merge_records(records: list[dict], *, fields: list[str] | None = None) -> MergedRecord:
"""Merge several ``{field: Candidate|value}`` records into one.
Each input record maps field name -> :class:`Candidate` (preferred) or a
bare value (treated as trust 0.5, no timestamp). The result keeps, for each
field, the winning value and the name of the source that supplied it.
"""
keys: list[str]
if fields is not None:
keys = list(fields)
else:
seen: dict[str, None] = {}
for rec in records:
for k in rec:
seen.setdefault(k, None)
keys = list(seen)
merged = MergedRecord()
for key in keys:
candidates: list[Candidate] = []
for rec in records:
if key not in rec:
continue
cand = rec[key]
if not isinstance(cand, Candidate):
cand = Candidate(value=cand, source="unknown")
candidates.append(cand)
resolution = resolve_field(candidates)
if resolution is not None:
merged.values[key] = resolution.value
if resolution.source is not None:
merged.provenance[key] = resolution.source
return merged
+176
View File
@@ -0,0 +1,176 @@
"""Product data-quality scoring.
The quality score is a 0..1 number combining four signals, per the locked
project decision:
quality = 0.4 * completeness
+ 0.3 * source_trust
+ 0.2 * multi_source_agreement
+ 0.1 * freshness
Each component is itself normalized to 0..1. The pure helpers below are
unit-testable; :func:`compute_quality` / :func:`update_quality` read the signals
for a product out of the database and persist the result on ``product``.
"""
from __future__ import annotations
from datetime import UTC, datetime
import psycopg
W_COMPLETENESS = 0.4
W_SOURCE_TRUST = 0.3
W_AGREEMENT = 0.2
W_FRESHNESS = 0.1
# Fields that count towards completeness (weighted equally).
COMPLETENESS_FIELDS = (
"name",
"gtin",
"brand",
"category",
"net_content",
"country_of_origin",
"nutriments",
"ingredients",
"image",
)
def completeness(present: set[str]) -> float:
"""Fraction of :data:`COMPLETENESS_FIELDS` that are present for a product."""
if not COMPLETENESS_FIELDS:
return 0.0
hits = sum(1 for f in COMPLETENESS_FIELDS if f in present)
return hits / len(COMPLETENESS_FIELDS)
def agreement_from_sources(source_count: int) -> float:
"""Multi-source corroboration proxy from the number of distinct sources.
A single source cannot be corroborated, so it scores a neutral 0.5; more
independent sources that describe the same product raise confidence.
"""
if source_count <= 1:
return 0.5
if source_count == 2:
return 0.8
return 1.0
def freshness_from_age(age_days: float | None) -> float:
"""Recency score from the age (in days) of the most recent source fetch."""
if age_days is None:
return 0.5
if age_days <= 30:
return 1.0
if age_days <= 180:
return 0.8
if age_days <= 365:
return 0.6
if age_days <= 730:
return 0.4
return 0.2
def score(
*,
completeness_score: float,
source_trust: float,
agreement: float,
freshness: float,
) -> float:
"""Combine the four normalized components into a 0..1 quality score."""
raw = (
W_COMPLETENESS * completeness_score
+ W_SOURCE_TRUST * source_trust
+ W_AGREEMENT * agreement
+ W_FRESHNESS * freshness
)
return round(max(0.0, min(1.0, raw)), 3)
def _present_fields(prod: dict, has_image: bool) -> set[str]:
present: set[str] = set()
if prod.get("name"):
present.add("name")
if prod.get("gtin"):
present.add("gtin")
if prod.get("brand_id"):
present.add("brand")
if prod.get("category_id"):
present.add("category")
if prod.get("net_content_canonical") is not None:
present.add("net_content")
if prod.get("country_of_origin"):
present.add("country_of_origin")
if prod.get("nutriments"):
present.add("nutriments")
if prod.get("ingredients_text"):
present.add("ingredients")
if has_image:
present.add("image")
return present
def compute_quality(conn: psycopg.Connection, product_id: str) -> float:
"""Compute (but do not persist) the quality score for one product."""
row = conn.execute(
"""
SELECT p.name, p.gtin, p.brand_id, p.category_id, p.net_content_canonical,
p.country_of_origin, f.nutriments, f.ingredients_text,
EXISTS (SELECT 1 FROM product_image pi WHERE pi.product_id = p.id)
FROM product p
LEFT JOIN food_detail f ON f.product_id = p.id
WHERE p.id = %s
""",
(product_id,),
).fetchone()
if row is None:
return 0.0
prod = {
"name": row[0],
"gtin": row[1],
"brand_id": row[2],
"category_id": row[3],
"net_content_canonical": row[4],
"country_of_origin": row[5],
"nutriments": row[6],
"ingredients_text": row[7],
}
has_image = bool(row[8])
src = conn.execute(
"""
SELECT count(DISTINCT ps.source_id), COALESCE(max(s.trust_weight), 0), max(ps.fetched_at)
FROM product_source ps
LEFT JOIN source s ON s.id = ps.source_id
WHERE ps.product_id = %s
""",
(product_id,),
).fetchone()
source_count = int(src[0] or 0)
source_trust = float(src[1] or 0.0)
last_fetched: datetime | None = src[2]
age_days: float | None = None
if last_fetched is not None:
now = datetime.now(UTC)
if last_fetched.tzinfo is None:
last_fetched = last_fetched.replace(tzinfo=UTC)
age_days = max(0.0, (now - last_fetched).total_seconds() / 86400.0)
return score(
completeness_score=completeness(_present_fields(prod, has_image)),
source_trust=source_trust,
agreement=agreement_from_sources(source_count),
freshness=freshness_from_age(age_days),
)
def update_quality(conn: psycopg.Connection, product_id: str) -> float:
"""Compute the quality score and write it to ``product.quality_score``."""
value = compute_quality(conn, product_id)
conn.execute("UPDATE product SET quality_score = %s WHERE id = %s", (value, product_id))
return value
+45
View File
@@ -0,0 +1,45 @@
"""Persistent ingestion watermark stored in the ``ingest_state`` table.
The incremental updater uses this to remember how far it got for each source
(e.g. Open Food Facts exposes a ``last_modified_t`` unix timestamp on every
product) so repeated runs only fetch what changed.
"""
from __future__ import annotations
from typing import Any
import psycopg
from psycopg.types.json import Jsonb
def get_watermark(conn: psycopg.Connection, source: str) -> int:
"""Return the last processed ``last_modified_t`` for *source* (0 if none)."""
row = conn.execute(
"SELECT last_modified_t FROM ingest_state WHERE source = %s", (source,)
).fetchone()
return int(row[0]) if row else 0
def set_watermark(
conn: psycopg.Connection,
source: str,
last_modified_t: int,
stats: dict[str, Any] | None = None,
) -> None:
"""Upsert the watermark and run metadata for *source*.
The watermark only ever moves forward: a lower ``last_modified_t`` is
ignored so an out-of-order or partial run cannot rewind progress.
"""
conn.execute(
"""
INSERT INTO ingest_state (source, last_modified_t, last_run_at, stats)
VALUES (%s, %s, now(), %s)
ON CONFLICT (source) DO UPDATE SET
last_modified_t = GREATEST(ingest_state.last_modified_t, EXCLUDED.last_modified_t),
last_run_at = now(),
stats = EXCLUDED.stats
""",
(source, int(last_modified_t), Jsonb(stats or {})),
)
+133
View File
@@ -0,0 +1,133 @@
"""Apply GS1 (or other authoritative) supplements to existing products.
A supplement only fills *gaps*: a field is written only when the product does
not already have a value. Each applied supplement records field-level provenance
in ``product_source`` and refreshes the product's quality score.
"""
from __future__ import annotations
from decimal import Decimal, InvalidOperation
from typing import Any
import psycopg
from psycopg.types.json import Jsonb
from opengoods import units
from opengoods.adapters.gs1 import GS1_HOMEPAGE, GS1_LICENSE, GS1_TRUST, SOURCE_NAME
from opengoods.etl.load import _ensure_brand, _normalize_brand, ensure_source_named
from opengoods.etl.quality import update_quality
def ensure_gs1_source(conn: psycopg.Connection) -> str:
"""Upsert the GS1 source row and return its id."""
return ensure_source_named(conn, SOURCE_NAME, GS1_HOMEPAGE, GS1_LICENSE, GS1_TRUST)
def _ensure_manufacturer(conn: psycopg.Connection, name: str | None) -> str | None:
if not name:
return None
row = conn.execute(
"""
INSERT INTO manufacturer (name, normalized_name)
VALUES (%s, %s)
ON CONFLICT (normalized_name) DO UPDATE SET name = manufacturer.name
RETURNING id
""",
(name, _normalize_brand(name)),
).fetchone()
return row[0]
def _net_content(rec: dict) -> tuple[Decimal, str, Decimal | None] | None:
raw_value = rec.get("net_content_value")
unit = rec.get("net_content_unit")
if raw_value is None or not unit:
return None
try:
value = Decimal(str(raw_value))
except (InvalidOperation, ValueError):
return None
try:
canonical = units.normalize(value, unit).canonical_value
except units.UnitError:
canonical = None
return value, unit, canonical
def apply_supplement(conn: psycopg.Connection, rec: dict[str, Any], source_id: str) -> list[str]:
"""Fill missing fields of the GTIN-matched product from ``rec``.
Returns the list of field names actually filled (empty if the product is
unknown or already complete for the supplied fields).
"""
gtin = rec.get("gtin")
if not gtin:
return []
prod = conn.execute(
"""
SELECT id, brand_id, manufacturer_id, gpc_brick_code, country_of_origin,
net_content_value
FROM product
WHERE gtin = %s AND status = 'active'
""",
(gtin,),
).fetchone()
if prod is None:
return []
product_id, brand_id, manufacturer_id, gpc, country, net_value = prod
sets: list[str] = []
params: list[Any] = []
filled: list[str] = []
if brand_id is None and rec.get("brand"):
new_brand_id = _ensure_brand(conn, rec["brand"])
if new_brand_id is not None:
sets.append("brand_id = %s")
params.append(new_brand_id)
filled.append("brand")
if manufacturer_id is None and rec.get("manufacturer"):
new_mfr_id = _ensure_manufacturer(conn, rec["manufacturer"])
if new_mfr_id is not None:
sets.append("manufacturer_id = %s")
params.append(new_mfr_id)
filled.append("manufacturer")
if gpc is None and rec.get("gpc_brick_code"):
sets.append("gpc_brick_code = %s")
params.append(rec["gpc_brick_code"])
filled.append("gpc_brick_code")
if country is None and rec.get("country_of_origin"):
sets.append("country_of_origin = %s")
params.append(rec["country_of_origin"])
filled.append("country_of_origin")
if net_value is None:
net = _net_content(rec)
if net is not None:
value, unit, canonical = net
sets += [
"net_content_value = %s",
"net_content_unit = %s",
"net_content_canonical = %s",
]
params += [value, unit, canonical]
filled.append("net_content")
if not filled:
return []
params.append(product_id)
conn.execute(f"UPDATE product SET {', '.join(sets)} WHERE id = %s", params)
conn.execute(
"""
INSERT INTO product_source (product_id, source_id, url, fields, fetched_at, raw)
VALUES (%s, %s, %s, %s, now(), %s)
""",
(product_id, source_id, GS1_HOMEPAGE, filled, Jsonb(rec)),
)
update_quality(conn, product_id)
return filled
+40
View File
@@ -0,0 +1,40 @@
"""Deduplicate products: merge non-GTIN duplicates into a canonical record.
Usage:
python -m opengoods.jobs.dedup --dry-run
python -m opengoods.jobs.dedup --actor nightly
"""
from __future__ import annotations
import argparse
import sys
import psycopg
from opengoods.etl.dedup import dedup_all
from opengoods.etl.load import default_dsn
def run(args: argparse.Namespace) -> int:
with psycopg.connect(args.dsn, autocommit=False) as conn:
summary = dedup_all(conn, actor=args.actor, dry_run=args.dry_run)
if args.dry_run:
conn.rollback()
else:
conn.commit()
mode = "dry-run" if args.dry_run else "applied"
print(f"{mode} groups={summary['groups']} merged={summary['merged']}")
return 0
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(description="Deduplicate OpenGoods products")
parser.add_argument("--actor", default="ingestion", help="merge_log actor label")
parser.add_argument("--dry-run", action="store_true", help="report only, do not write")
parser.add_argument("--dsn", default=default_dsn(), help="PostgreSQL DSN")
return run(parser.parse_args(argv))
if __name__ == "__main__":
sys.exit(main())
+66
View File
@@ -0,0 +1,66 @@
"""Lightweight recurring ingestion scheduler.
Runs one ingestion cycle (incremental OFF update, then dedup) on a fixed
interval. Dependency-free: a plain sleep loop rather than a cron/APScheduler
dependency, so it is trivial to run in a container or under systemd/supervisor.
Usage:
python -m opengoods.jobs.schedule --once # single cycle, then exit
python -m opengoods.jobs.schedule --interval 3600 # every hour
"""
from __future__ import annotations
import argparse
import sys
import time
from datetime import UTC, datetime
from opengoods.etl.load import default_dsn
from opengoods.jobs import dedup as dedup_job
from opengoods.jobs import update_off as update_job
def _cycle(args: argparse.Namespace) -> None:
ts = datetime.now(UTC).isoformat(timespec="seconds")
print(f"[{ts}] cycle start")
update_job.run(
argparse.Namespace(
since=None,
page_size=args.page_size,
max_pages=args.max_pages,
min_interval=args.min_interval,
dsn=args.dsn,
)
)
if not args.skip_dedup:
dedup_job.run(argparse.Namespace(actor="scheduler", dry_run=False, dsn=args.dsn))
def run(args: argparse.Namespace) -> int:
_cycle(args)
if args.once:
return 0
while True:
time.sleep(args.interval)
try:
_cycle(args)
except Exception as exc: # noqa: BLE001 - keep the loop alive across failures
print(f"cycle error: {exc}", file=sys.stderr)
return 0
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(description="Recurring OpenGoods ingestion")
parser.add_argument("--interval", type=int, default=3600, help="seconds between cycles")
parser.add_argument("--once", action="store_true", help="run a single cycle and exit")
parser.add_argument("--skip-dedup", action="store_true", help="run update only")
parser.add_argument("--page-size", type=int, default=100, help="search page size")
parser.add_argument("--max-pages", type=int, default=10, help="max pages to scan")
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))
if __name__ == "__main__":
sys.exit(main())
+65
View File
@@ -0,0 +1,65 @@
"""Incremental Open Food Facts update.
Fetches products modified since the persisted watermark, loads them, then
advances the watermark to the newest ``last_modified_t`` processed so the next
run only sees what changed.
Usage:
python -m opengoods.jobs.update_off --max-pages 5
python -m opengoods.jobs.update_off --since 1700000000 # override watermark
"""
from __future__ import annotations
import argparse
import sys
import psycopg
from opengoods.adapters.openfoodfacts import SOURCE_NAME, OpenFoodFactsAdapter
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 = 0
high_watermark = 0
with psycopg.connect(args.dsn, autocommit=False) as conn:
source_id = ensure_source(conn)
since = args.since if args.since is not None else get_watermark(conn, SOURCE_NAME)
high_watermark = since
for raw in adapter.fetch_modified_since(
since, page_size=args.page_size, max_pages=args.max_pages
):
high_watermark = max(high_watermark, int(raw.get("last_modified_t") or 0))
rec = transform(raw)
if rec is None:
skipped += 1
continue
load_record(conn, rec, source_id, raw)
loaded += 1
set_watermark(
conn,
SOURCE_NAME,
high_watermark,
stats={"loaded": loaded, "skipped": skipped, "since": since},
)
conn.commit()
print(f"since={since} loaded={loaded} skipped={skipped} watermark={high_watermark}")
return 0
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(description="Incremental OFF update")
parser.add_argument("--since", type=int, default=None, help="override watermark (unix ts)")
parser.add_argument("--page-size", type=int, default=100, help="search page size")
parser.add_argument("--max-pages", type=int, default=10, help="max pages to scan")
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))
if __name__ == "__main__":
sys.exit(main())
+31
View File
@@ -0,0 +1,31 @@
"""Shared test fixtures.
`db_conn` yields a psycopg connection inside a transaction that is rolled back
after each test, so DB tests stay isolated and leave no residue. Tests are
skipped automatically when no database is reachable or M4 migrations are not
applied (e.g. local runs without docker).
"""
from __future__ import annotations
import psycopg
import pytest
from opengoods.etl.load import default_dsn
@pytest.fixture()
def db_conn():
try:
conn = psycopg.connect(default_dsn(), connect_timeout=3)
except psycopg.OperationalError as exc: # pragma: no cover - env dependent
pytest.skip(f"no database available: {exc}")
has_state = conn.execute("SELECT to_regclass('public.ingest_state') IS NOT NULL").fetchone()[0]
if not has_state:
conn.close()
pytest.skip("M4 migrations not applied")
try:
yield conn
finally:
conn.rollback()
conn.close()
+14
View File
@@ -0,0 +1,14 @@
{
"06901234567892": {
"name": "示例矿泉水 550ml",
"brand": "示例品牌",
"manufacturer": "示例饮品有限公司",
"gpc_brick_code": "10000224",
"country_of_origin": "China",
"net_content_value": 550,
"net_content_unit": "ml"
},
"00000000000000": {
"brand": ""
}
}
+62
View File
@@ -0,0 +1,62 @@
from opengoods.etl.dedup import choose_canonical, dedup_all, product_signature
from opengoods.etl.load import _ensure_brand, ensure_source
def test_product_signature_normalization():
a = product_signature(" Spring Water ", "Acme", 500)
b = product_signature("spring water", "acme", 500)
assert a == b
assert product_signature("", "x", 1) is None
def test_choose_canonical_prefers_quality():
members = [
{"id": "a", "quality_score": 0.2, "created_at": 1},
{"id": "b", "quality_score": 0.9, "created_at": 2},
]
assert choose_canonical(members)["id"] == "b"
def test_dedup_merges_duplicates(db_conn):
brand_id = _ensure_brand(db_conn, "DupBrand")
def mk(quality):
return db_conn.execute(
"""
INSERT INTO product (name, brand_id, net_content_canonical, quality_score)
VALUES (%s, %s, %s, %s) RETURNING id
""",
("Dup Snack", brand_id, 100, quality),
).fetchone()[0]
keep = mk(0.9)
drop = mk(0.2)
src = ensure_source(db_conn)
db_conn.execute(
"INSERT INTO product_source (product_id, source_id, fields) VALUES (%s, %s, %s)",
(drop, src, ["name"]),
)
summary = dedup_all(db_conn)
assert summary == {"groups": 1, "merged": 1}
keep_status = db_conn.execute("SELECT status FROM product WHERE id = %s", (keep,)).fetchone()[0]
drop_status, canonical_id = db_conn.execute(
"SELECT status, canonical_id FROM product WHERE id = %s", (drop,)
).fetchone()
assert keep_status == "active"
assert drop_status == "merged"
assert str(canonical_id) == str(keep)
# The merged product's source row was re-pointed to the canonical product.
reattached = db_conn.execute(
"SELECT count(*) FROM product_source WHERE product_id = %s", (keep,)
).fetchone()[0]
assert reattached == 1
logged = db_conn.execute(
"SELECT count(*) FROM merge_log WHERE kept_id = %s AND merged_id = %s",
(keep, drop),
).fetchone()[0]
assert logged == 1
+53
View File
@@ -0,0 +1,53 @@
from pathlib import Path
from opengoods.adapters.gs1 import GS1Adapter
from opengoods.etl.supplement import apply_supplement, ensure_gs1_source
MAPPING = Path(__file__).parent / "fixtures" / "gs1_mapping.json"
GTIN = "06901234567892"
def test_gs1_adapter_offline_lookup():
adapter = GS1Adapter.from_file(MAPPING)
rec = adapter.fetch_barcode(GTIN)
assert rec["brand"] == "示例品牌"
assert rec["net_content_value"] == 550
assert rec["net_content_unit"] == "ml"
# An entry that only has empty values yields no supplement.
assert adapter.fetch_barcode("00000000000000") is None
# Unknown barcode -> None.
assert adapter.fetch_barcode("99999999999999") is None
def test_gs1_supplement_fills_only_gaps(db_conn):
pid = db_conn.execute(
"INSERT INTO product (gtin, name) VALUES (%s, %s) RETURNING id", (GTIN, "")
).fetchone()[0]
adapter = GS1Adapter.from_file(MAPPING)
rec = adapter.fetch_barcode(GTIN)
source_id = ensure_gs1_source(db_conn)
filled = apply_supplement(db_conn, rec, source_id)
assert {"brand", "country_of_origin", "net_content"} <= set(filled)
brand_id, country, net_value, net_unit = db_conn.execute(
"""
SELECT brand_id, country_of_origin, net_content_value, net_content_unit
FROM product WHERE id = %s
""",
(pid,),
).fetchone()
assert brand_id is not None
assert country == "China"
assert float(net_value) == 550.0
assert net_unit == "ml"
fields = db_conn.execute(
"SELECT fields FROM product_source WHERE product_id = %s AND source_id = %s",
(pid, source_id),
).fetchone()[0]
assert "brand" in fields
# Re-applying does nothing because the gaps are now filled.
assert apply_supplement(db_conn, rec, source_id) == []
+62
View File
@@ -0,0 +1,62 @@
from datetime import UTC, datetime
from opengoods.etl.merge import Candidate, merge_records, resolve_field
def _ts(y, m, d):
return datetime(y, m, d, tzinfo=UTC)
def test_resolve_field_prefers_trust_then_recency():
cands = [
Candidate(value="A", source="off", trust=0.7, fetched_at=_ts(2024, 1, 1)),
Candidate(value="B", source="gs1", trust=0.9, fetched_at=_ts(2023, 1, 1)),
]
res = resolve_field(cands)
assert res is not None
assert res.value == "B"
assert res.source == "gs1"
def test_resolve_field_recency_tiebreak_on_equal_trust():
cands = [
Candidate(value="old", source="a", trust=0.7, fetched_at=_ts(2023, 1, 1)),
Candidate(value="new", source="b", trust=0.7, fetched_at=_ts(2024, 6, 1)),
]
assert resolve_field(cands).value == "new"
def test_resolve_field_skips_empty():
cands = [
Candidate(value="", source="a", trust=0.99),
Candidate(value=None, source="b", trust=0.99),
Candidate(value="kept", source="c", trust=0.1),
]
assert resolve_field(cands).value == "kept"
assert resolve_field([Candidate(value="", source="a")]) is None
def test_merge_records_provenance():
records = [
{
"name": Candidate("Water", "off", 0.7, _ts(2024, 1, 1)),
"brand": Candidate("", "off", 0.7),
},
{
"brand": Candidate("Acme", "gs1", 0.9, _ts(2024, 2, 1)),
"gtin": Candidate("123", "gs1", 0.9),
},
]
merged = merge_records(records)
assert merged.values["name"] == "Water"
assert merged.values["brand"] == "Acme"
assert merged.values["gtin"] == "123"
assert merged.provenance["brand"] == "gs1"
assert merged.provenance["name"] == "off"
def test_merge_records_accepts_bare_values():
merged = merge_records([{"x": 1}, {"x": 2}])
# both bare -> trust tie, no timestamps -> first max() wins deterministically
assert merged.values["x"] in (1, 2)
assert merged.provenance["x"] == "unknown"
+44
View File
@@ -0,0 +1,44 @@
import httpx
from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter
def _product(code, lm):
return {"code": code, "product_name": f"P{code}", "last_modified_t": lm}
def _adapter(pages):
"""Build an adapter whose search endpoint serves the given pages."""
def handler(request: httpx.Request) -> httpx.Response:
page = int(request.url.params.get("page", "1"))
products = pages.get(page, [])
return httpx.Response(200, json={"products": products, "page": page})
client = httpx.Client(transport=httpx.MockTransport(handler))
return OpenFoodFactsAdapter(client=client, min_interval=0)
def test_incremental_yields_only_newer_and_stops_at_watermark():
pages = {
1: [_product("1", 300), _product("2", 250), _product("3", 100)],
}
adapter = _adapter(pages)
got = list(adapter.fetch_modified_since(200, page_size=3, max_pages=5))
codes = [p["code"] for p in got]
assert codes == ["1", "2"] # 100 <= 200 stops iteration
def test_incremental_paginates_until_short_page():
pages = {
1: [_product("1", 900), _product("2", 800)],
2: [_product("3", 700)], # short page -> stop after
}
adapter = _adapter(pages)
got = list(adapter.fetch_modified_since(0, page_size=2, max_pages=5))
assert [p["code"] for p in got] == ["1", "2", "3"]
def test_incremental_empty_first_page():
adapter = _adapter({1: []})
assert list(adapter.fetch_modified_since(0, page_size=10, max_pages=3)) == []
+38
View File
@@ -0,0 +1,38 @@
from opengoods.etl.quality import (
COMPLETENESS_FIELDS,
agreement_from_sources,
completeness,
freshness_from_age,
score,
)
def test_completeness_bounds():
assert completeness(set()) == 0.0
assert completeness(set(COMPLETENESS_FIELDS)) == 1.0
half = set(list(COMPLETENESS_FIELDS)[: len(COMPLETENESS_FIELDS) // 2])
assert 0.0 < completeness(half) < 1.0
def test_agreement_from_sources():
assert agreement_from_sources(0) == 0.5
assert agreement_from_sources(1) == 0.5
assert agreement_from_sources(2) == 0.8
assert agreement_from_sources(5) == 1.0
def test_freshness_from_age():
assert freshness_from_age(None) == 0.5
assert freshness_from_age(1) == 1.0
assert freshness_from_age(100) == 0.8
assert freshness_from_age(300) == 0.6
assert freshness_from_age(700) == 0.4
assert freshness_from_age(5000) == 0.2
def test_score_weighted_sum_and_bounds():
assert score(completeness_score=0, source_trust=0, agreement=0, freshness=0) == 0.0
assert score(completeness_score=1, source_trust=1, agreement=1, freshness=1) == 1.0
# 0.4*1 + 0.3*0.5 + 0.2*0.5 + 0.1*1 = 0.75
got = score(completeness_score=1.0, source_trust=0.5, agreement=0.5, freshness=1.0)
assert got == 0.75
+21
View File
@@ -0,0 +1,21 @@
import json
from pathlib import Path
from opengoods.etl.load import ensure_source, load_record
from opengoods.etl.quality import compute_quality
from opengoods.etl.transform import transform
FIXTURE = json.loads((Path(__file__).parent / "fixtures" / "off_product.json").read_text())
def test_quality_score_set_on_load(db_conn):
source_id = ensure_source(db_conn)
rec = transform(FIXTURE)
pid = load_record(db_conn, rec, source_id, FIXTURE)
stored = float(
db_conn.execute("SELECT quality_score FROM product WHERE id = %s", (pid,)).fetchone()[0]
)
assert 0.0 < stored <= 1.0
# The persisted value matches a fresh recomputation.
assert abs(stored - compute_quality(db_conn, pid)) < 1e-9
+16
View File
@@ -0,0 +1,16 @@
from opengoods.etl.state import get_watermark, set_watermark
def test_watermark_roundtrip_and_monotonic(db_conn):
src = "test-source"
assert get_watermark(db_conn, src) == 0
set_watermark(db_conn, src, 100, stats={"loaded": 1})
assert get_watermark(db_conn, src) == 100
# A lower watermark must not rewind progress.
set_watermark(db_conn, src, 50)
assert get_watermark(db_conn, src) == 100
set_watermark(db_conn, src, 150)
assert get_watermark(db_conn, src) == 150
+1
View File
@@ -0,0 +1 @@
DROP TABLE IF EXISTS ingest_state;
+10
View File
@@ -0,0 +1,10 @@
-- M4 ingestion management: persistent per-source incremental watermark.
-- The updater reads/writes one row per source to resume incremental imports
-- (e.g. Open Food Facts `last_modified_t`) and to record run statistics.
CREATE TABLE ingest_state (
source TEXT PRIMARY KEY,
last_modified_t BIGINT NOT NULL DEFAULT 0,
last_run_at TIMESTAMPTZ,
cursor TEXT,
stats JSONB NOT NULL DEFAULT '{}'
);
+4
View File
@@ -0,0 +1,4 @@
DROP INDEX IF EXISTS idx_audit_created;
DROP INDEX IF EXISTS idx_audit_entity;
DROP TABLE IF EXISTS audit_log;
DELETE FROM source WHERE name = 'manual';
+22
View File
@@ -0,0 +1,22 @@
-- Admin console support: manual edit provenance + audit log.
-- Manual-entry source used to record field-level provenance for operator edits.
INSERT INTO source (name, homepage, license, trust_weight, notes)
VALUES ('manual', NULL, 'proprietary', 0.90, '运营人工录入/补全')
ON CONFLICT (name) DO NOTHING;
-- Audit trail: every admin write records actor, action, entity and before/after.
CREATE TABLE IF NOT EXISTS audit_log (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
actor TEXT NOT NULL,
action TEXT NOT NULL,
entity TEXT NOT NULL,
entity_id UUID,
fields TEXT[] NOT NULL DEFAULT '{}',
before JSONB,
after JSONB,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX IF NOT EXISTS idx_audit_entity ON audit_log (entity, entity_id);
CREATE INDEX IF NOT EXISTS idx_audit_created ON audit_log (created_at DESC);
+1
View File
@@ -0,0 +1 @@
DROP TABLE IF EXISTS submission;
+30
View File
@@ -0,0 +1,30 @@
-- Community contributions: public users may submit new/supplementary product
-- archives. Submissions go to a moderation queue and never touch the product
-- tables until an admin approves them.
-- Community source used to record field-level provenance for approved
-- public contributions (lower trust than manual operator edits).
INSERT INTO source (name, homepage, license, trust_weight, notes)
VALUES ('community', NULL, 'user-contributed', 0.50, '公众投稿/众包贡献,经人工审核后收纳')
ON CONFLICT (name) DO NOTHING;
CREATE TABLE IF NOT EXISTS submission (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
gtin VARCHAR(14),
name TEXT NOT NULL,
payload JSONB NOT NULL DEFAULT '{}',
target_product_id UUID REFERENCES product(id) ON DELETE SET NULL,
result_product_id UUID REFERENCES product(id) ON DELETE SET NULL,
submitter_name TEXT,
submitter_contact TEXT,
note TEXT,
status VARCHAR(16) NOT NULL DEFAULT 'pending',
review_note TEXT,
reviewed_by TEXT,
reviewed_at TIMESTAMPTZ,
remote_ip TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
CONSTRAINT submission_status_chk CHECK (status IN ('pending','approved','rejected'))
);
CREATE INDEX IF NOT EXISTS idx_submission_status ON submission (status, created_at DESC);
+12
View File
@@ -0,0 +1,12 @@
<!doctype html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8" />
<meta name="viewport" content="width=device-width, initial-scale=1.0" />
<title>天工商品档案公共仓</title>
</head>
<body>
<div id="root"></div>
<script type="module" src="/src/main.tsx"></script>
</body>
</html>
+2690
View File
File diff suppressed because it is too large Load Diff
+26
View File
@@ -0,0 +1,26 @@
{
"name": "opengoods-public-frontend",
"private": true,
"version": "1.0.0",
"type": "module",
"scripts": {
"dev": "vite",
"build": "tsc && vite build",
"preview": "vite preview"
},
"dependencies": {
"react": "^18.2.0",
"react-dom": "^18.2.0",
"lucide-react": "^0.344.0"
},
"devDependencies": {
"@types/react": "^18.2.55",
"@types/react-dom": "^18.2.19",
"@vitejs/plugin-react": "^4.2.1",
"autoprefixer": "^10.4.17",
"postcss": "^8.4.35",
"tailwindcss": "^3.4.1",
"typescript": "^5.3.3",
"vite": "^5.1.0"
}
}
+6
View File
@@ -0,0 +1,6 @@
export default {
plugins: {
tailwindcss: {},
autoprefixer: {},
},
};
+95
View File
@@ -0,0 +1,95 @@
import { useState } from "react";
import { Boxes, Search, PlusCircle, Code2 } from "lucide-react";
import Home from "./components/Home";
import ProductView from "./components/ProductView";
import Contribute from "./components/Contribute";
import ApiDocs from "./components/ApiDocs";
type View =
| { name: "home" }
| { name: "product"; id: string }
| { name: "contribute" }
| { name: "api" };
export default function App() {
const [view, setView] = useState<View>({ name: "home" });
return (
<div className="min-h-full flex flex-col">
<header className="bg-white border-b">
<div className="max-w-5xl mx-auto px-4 h-14 flex items-center justify-between">
<button
className="flex items-center gap-2 font-semibold text-gray-800"
onClick={() => setView({ name: "home" })}
>
<Boxes className="w-6 h-6 text-emerald-600" />
<span className="text-gray-400 font-normal text-sm"></span>
</button>
<nav className="flex items-center gap-1 text-sm">
<button
className={`px-3 py-1.5 rounded-md flex items-center gap-1.5 ${
view.name === "home" ? "bg-emerald-50 text-emerald-700" : "text-gray-600 hover:bg-gray-100"
}`}
onClick={() => setView({ name: "home" })}
>
<Search className="w-4 h-4" />
</button>
<button
className={`px-3 py-1.5 rounded-md flex items-center gap-1.5 ${
view.name === "contribute" ? "bg-emerald-50 text-emerald-700" : "text-gray-600 hover:bg-gray-100"
}`}
onClick={() => setView({ name: "contribute" })}
>
<PlusCircle className="w-4 h-4" />
</button>
<button
className={`px-3 py-1.5 rounded-md flex items-center gap-1.5 ${
view.name === "api" ? "bg-emerald-50 text-emerald-700" : "text-gray-600 hover:bg-gray-100"
}`}
onClick={() => setView({ name: "api" })}
>
<Code2 className="w-4 h-4" /> API
</button>
</nav>
</div>
</header>
<main className="flex-1 max-w-5xl w-full mx-auto px-4 py-6">
{view.name === "home" && (
<Home
onOpen={(id) => setView({ name: "product", id })}
onContribute={() => setView({ name: "contribute" })}
onApi={() => setView({ name: "api" })}
/>
)}
{view.name === "product" && (
<ProductView id={view.id} onBack={() => setView({ name: "home" })} />
)}
{view.name === "contribute" && (
<Contribute onDone={() => setView({ name: "home" })} />
)}
{view.name === "api" && <ApiDocs />}
</main>
<footer className="border-t bg-white">
<div className="max-w-5xl mx-auto px-4 py-4 text-xs text-gray-400 leading-relaxed">
/
稿
<button onClick={() => setView({ name: "api" })} className="ml-1 text-emerald-600 hover:underline">
API
</button>
<div className="mt-2">
<a
href="https://beian.miit.gov.cn/"
target="_blank"
rel="noreferrer"
className="hover:text-gray-600"
>
ICP备2025185218号-6
</a>
</div>
</div>
</footer>
</div>
);
}
+40
View File
@@ -0,0 +1,40 @@
import type { Category, Product, ProductSummary, SubmissionInput } from "./types";
async function req<T>(path: string, init?: RequestInit): Promise<T> {
const res = await fetch(path, {
...init,
headers: { "Content-Type": "application/json", ...(init?.headers || {}) },
});
if (!res.ok) {
let msg = `请求失败 (${res.status})`;
try {
const body = await res.json();
if (body?.error?.message) msg = body.error.message;
} catch {
/* ignore */
}
throw new Error(msg);
}
return res.json() as Promise<T>;
}
export interface SearchResult {
items: ProductSummary[];
page: number;
size: number;
total: number;
}
export const api = {
search: (q: string, page = 1, size = 20) =>
req<SearchResult>(
`/api/v1/products/search?q=${encodeURIComponent(q)}&page=${page}&size=${size}`,
),
product: (id: string) => req<Product>(`/api/v1/products/${id}`),
categories: () => req<{ items: Category[] }>(`/api/v1/categories`),
submit: (input: SubmissionInput) =>
req<{ id: string; status: string }>(`/api/public/submissions`, {
method: "POST",
body: JSON.stringify(input),
}),
};
+279
View File
@@ -0,0 +1,279 @@
import { useState } from "react";
import { Check, Copy } from "lucide-react";
const ORIGIN = typeof window !== "undefined" ? window.location.origin : "https://goods.tangshasha.com";
const BASE = `${ORIGIN}/api/v1`;
function CopyBtn({ text }: { text: string }) {
const [done, setDone] = useState(false);
return (
<button
onClick={async () => {
try {
await navigator.clipboard.writeText(text);
setDone(true);
setTimeout(() => setDone(false), 1200);
} catch {
/* clipboard unavailable */
}
}}
className="text-gray-400 hover:text-gray-600"
title="复制"
>
{done ? <Check className="w-4 h-4 text-emerald-600" /> : <Copy className="w-4 h-4" />}
</button>
);
}
function Code({ children }: { children: string }) {
return (
<div className="relative group">
<pre className="bg-gray-900 text-gray-100 text-xs rounded-md p-3 overflow-x-auto whitespace-pre">
{children}
</pre>
<div className="absolute top-2 right-2 opacity-70 group-hover:opacity-100">
<CopyBtn text={children} />
</div>
</div>
);
}
function Method({ m }: { m: string }) {
const color = m === "GET" ? "bg-sky-100 text-sky-700" : "bg-emerald-100 text-emerald-700";
return <span className={`text-xs font-mono font-semibold rounded px-1.5 py-0.5 ${color}`}>{m}</span>;
}
type Param = { name: string; required?: boolean; desc: string };
function Endpoint({
method,
path,
title,
desc,
params,
example,
response,
}: {
method: string;
path: string;
title: string;
desc: string;
params?: Param[];
example: string;
response: string;
}) {
return (
<div className="bg-white border rounded-lg p-5">
<div className="flex items-center gap-2 flex-wrap">
<Method m={method} />
<code className="text-sm text-gray-800 font-mono break-all">{path}</code>
<span className="ml-auto" />
<CopyBtn text={`${ORIGIN}${path}`} />
</div>
<div className="mt-2 font-medium text-gray-800">{title}</div>
<p className="text-sm text-gray-500 mt-0.5">{desc}</p>
{params && params.length > 0 && (
<table className="mt-3 w-full text-sm">
<thead className="text-gray-400 text-left">
<tr>
<th className="font-medium pr-4 pb-1"></th>
<th className="font-medium pr-4 pb-1"></th>
<th className="font-medium pb-1"></th>
</tr>
</thead>
<tbody className="align-top">
{params.map((p) => (
<tr key={p.name}>
<td className="pr-4 py-0.5 font-mono text-gray-700">{p.name}</td>
<td className="pr-4 py-0.5 text-gray-500">{p.required ? "是" : "否"}</td>
<td className="py-0.5 text-gray-600">{p.desc}</td>
</tr>
))}
</tbody>
</table>
)}
<div className="mt-3 text-xs text-gray-400 mb-1"></div>
<Code>{example}</Code>
<div className="mt-3 text-xs text-gray-400 mb-1"></div>
<Code>{response}</Code>
</div>
);
}
export default function ApiDocs() {
return (
<div className="space-y-5">
<div className="bg-white border rounded-lg p-5">
<h1 className="text-2xl font-bold text-gray-800">API </h1>
<p className="mt-2 text-gray-600 text-sm leading-relaxed">
<strong></strong> REST API
/Nutri-Score
JSONUTF-8/
</p>
<div className="mt-3 text-sm text-gray-700">
<div>
<code className="font-mono bg-gray-100 rounded px-1.5 py-0.5">{BASE}</code>
</div>
<ul className="mt-2 list-disc pl-5 text-gray-600 space-y-1">
<li> API Key / Token GET </li>
<li>
<code className="font-mono">page</code> 1
<code className="font-mono">size</code> 20 100
</li>
<li>
<code className="font-mono">404</code>
<code className="font-mono">{` {"error":{"code","message","request_id"}}`}</code>
</li>
<li></li>
</ul>
</div>
</div>
<Endpoint
method="GET"
path="/healthz"
title="健康检查"
desc="服务存活探针。"
example={`curl ${ORIGIN}/healthz`}
response={`{ "status": "ok" }`}
/>
<Endpoint
method="GET"
path="/api/v1/products/barcode/{gtin}"
title="按条码查询商品"
desc="按 GTIN(条码)精确查询单个商品档案。"
params={[{ name: "gtin", required: true, desc: "商品条码(路径参数),如 5449000000996" }]}
example={`curl ${BASE}/products/barcode/5449000000996`}
response={`{
"id": "…",
"gtin": "5449000000996",
"name": "可口可乐 经典原味",
"brand": "Coca-Cola",
"category_path": "food.beverages.carbonated",
"net_content_value": 33,
"net_content_unit": "cl",
"country_of_origin": "Algeria",
"quality_score": 0.81
}`}
/>
<Endpoint
method="GET"
path="/api/v1/products/search"
title="搜索商品"
desc="按名称模糊搜索,可按品类过滤,支持分页。"
params={[
{ name: "q", desc: "关键词(名称/条码),留空返回全部" },
{ name: "category", desc: "品类编码过滤,如 food.beverages" },
{ name: "page", desc: "页码,默认 1" },
{ name: "size", desc: "每页条数,默认 20,最大 100" },
]}
example={`curl "${BASE}/products/search?q=nutella&page=1&size=20"`}
response={`{
"items": [
{ "id": "…", "gtin": "3017624010701",
"name": "Nutella", "brand": "Ferrero",
"category_path": "food.snacks.chocolate" }
],
"page": 1, "size": 20, "total": 1
}`}
/>
<Endpoint
method="GET"
path="/api/v1/products/{id}"
title="商品详情"
desc="按商品 UUID 获取完整档案(含配料、营养、添加剂、图片、MSRP 等)。"
params={[{ name: "id", required: true, desc: "商品 UUID(路径参数)" }]}
example={`curl ${BASE}/products/{id}`}
response={`{
"id": "…", "name": "…", "brand": "…",
"ingredients_text": "…",
"nutriments": { "energy_kcal": 42, "sugars_g": 10.6 },
"nutrition_basis": "per_100g",
"nutri_score": "E",
"images": [], "msrp": [],
"quality_score": 0.81
}`}
/>
<Endpoint
method="GET"
path="/api/v1/products/{id}/nutriments"
title="商品营养成分"
desc="仅返回该商品的营养字段。"
params={[{ name: "id", required: true, desc: "商品 UUID(路径参数)" }]}
example={`curl ${BASE}/products/{id}/nutriments`}
response={`{
"basis": "per_100g",
"values": { "energy_kcal": 42, "sugars_g": 10.6, "salt_g": 0 }
}`}
/>
<Endpoint
method="GET"
path="/api/v1/products/{id}/msrp"
title="厂商建议零售价快照"
desc="官方建议零售价历史快照,仅供参考,不含任何购买入口。"
params={[{ name: "id", required: true, desc: "商品 UUID(路径参数)" }]}
example={`curl ${BASE}/products/{id}/msrp`}
response={`{
"items": [
{ "amount": 3.5, "currency": "CNY", "region": "CN", "effective_date": "2025-01-01" }
],
"disclaimer": "厂商建议零售价历史快照,仅供参考,不构成购买建议…"
}`}
/>
<Endpoint
method="GET"
path="/api/v1/brands"
title="品牌列表"
desc="分页列出全部品牌。"
params={[
{ name: "page", desc: "页码,默认 1" },
{ name: "size", desc: "每页条数,默认 20,最大 100" },
]}
example={`curl "${BASE}/brands?page=1&size=20"`}
response={`{
"items": [ { "id": "…", "name": "Ferrero" } ],
"page": 1, "size": 20, "total": 4
}`}
/>
<Endpoint
method="GET"
path="/api/v1/categories"
title="品类树"
desc="返回完整品类层级(编码 + 名称)。"
example={`curl ${BASE}/categories`}
response={`{
"items": [
{ "id": "…", "code": "food", "name": "食品饮料", "parent_id": null }
]
}`}
/>
<Endpoint
method="GET"
path="/api/v1/sources/{id}"
title="数据来源"
desc="按 ID 查询单个数据来源(含许可、信任权重)。"
params={[{ name: "id", required: true, desc: "来源 UUID(路径参数)" }]}
example={`curl ${BASE}/sources/{id}`}
response={`{
"id": "…", "name": "openfoodfacts",
"license": "ODbL", "trust_weight": 0.7
}`}
/>
<div className="text-xs text-gray-400 leading-relaxed">
/
OpenFoodFacts ODbL
</div>
</div>
);
}
@@ -0,0 +1,276 @@
import { useEffect, useState } from "react";
import { CheckCircle2, PlusCircle, Trash2 } from "lucide-react";
import { api } from "../api";
import type { Category, SubmissionImage, SubmissionInput } from "../types";
const NUTRI_FIELDS: { key: string; label: string }[] = [
{ key: "energy_kcal", label: "能量 (kcal)" },
{ key: "energy_kj", label: "能量 (kJ)" },
{ key: "fat", label: "脂肪 (g)" },
{ key: "saturated_fat", label: "饱和脂肪 (g)" },
{ key: "carbohydrates", label: "碳水 (g)" },
{ key: "sugars", label: "糖 (g)" },
{ key: "proteins", label: "蛋白质 (g)" },
{ key: "salt", label: "盐 (g)" },
];
function field(v: string): string | null {
const t = v.trim();
return t === "" ? null : t;
}
export default function Contribute({ onDone }: { onDone: () => void }) {
const [categories, setCategories] = useState<Category[]>([]);
const [done, setDone] = useState(false);
const [submitting, setSubmitting] = useState(false);
const [error, setError] = useState("");
const [name, setName] = useState("");
const [gtin, setGtin] = useState("");
const [brand, setBrand] = useState("");
const [categoryID, setCategoryID] = useState("");
const [netValue, setNetValue] = useState("");
const [netUnit, setNetUnit] = useState("");
const [country, setCountry] = useState("");
const [ingredients, setIngredients] = useState("");
const [basis, setBasis] = useState("");
const [nutri, setNutri] = useState<Record<string, string>>({});
const [images, setImages] = useState<SubmissionImage[]>([]);
const [imageURL, setImageURL] = useState("");
const [submitter, setSubmitter] = useState("");
const [contact, setContact] = useState("");
const [note, setNote] = useState("");
useEffect(() => {
api.categories().then((r) => setCategories(r.items)).catch(() => undefined);
}, []);
async function submit(e: React.FormEvent) {
e.preventDefault();
if (name.trim() === "") {
setError("请填写商品名称");
return;
}
setSubmitting(true);
setError("");
const nutriments: Record<string, number> = {};
for (const [k, v] of Object.entries(nutri)) {
const n = parseFloat(v);
if (!Number.isNaN(n)) nutriments[k] = n;
}
const input: SubmissionInput = {
name: name.trim(),
gtin: field(gtin),
brand_name: field(brand),
category_id: categoryID || null,
net_content_value: field(netValue) ? parseFloat(netValue) : null,
net_content_unit: field(netUnit),
country_of_origin: field(country),
ingredients_text: field(ingredients),
nutriments: Object.keys(nutriments).length ? nutriments : null,
nutrition_basis: basis || null,
images: images.length ? images : undefined,
submitter_name: field(submitter),
submitter_contact: field(contact),
note: field(note),
};
try {
await api.submit(input);
setDone(true);
} catch (err) {
setError(err instanceof Error ? err.message : "提交失败");
} finally {
setSubmitting(false);
}
}
if (done) {
return (
<div className="max-w-xl mx-auto text-center py-16">
<CheckCircle2 className="w-14 h-14 text-emerald-500 mx-auto" />
<h1 className="mt-4 text-xl font-semibold text-gray-800"></h1>
<p className="mt-2 text-gray-500">
</p>
<button
onClick={onDone}
className="mt-6 px-5 py-2 rounded-lg bg-emerald-600 text-white hover:bg-emerald-700"
>
</button>
</div>
);
}
const input =
"w-full border rounded-md px-3 py-2 text-sm focus:outline-none focus:ring-2 focus:ring-emerald-400";
const label = "block text-xs text-gray-500 mb-1";
return (
<form onSubmit={submit} className="max-w-3xl mx-auto">
<h1 className="text-xl font-semibold text-gray-800"></h1>
<p className="mt-1 text-sm text-gray-500">
<b></b> *
</p>
{error && (
<div className="mt-4 bg-red-50 text-red-700 text-sm rounded-md px-4 py-2">{error}</div>
)}
<div className="bg-white border rounded-lg p-5 mt-4">
<h2 className="font-medium text-gray-700 mb-3"></h2>
<div className="grid grid-cols-1 sm:grid-cols-2 gap-4">
<div className="sm:col-span-2">
<label className={label}> *</label>
<input className={input} value={name} onChange={(e) => setName(e.target.value)} />
</div>
<div>
<label className={label}> (GTIN)</label>
<input className={input} value={gtin} onChange={(e) => setGtin(e.target.value)} />
</div>
<div>
<label className={label}></label>
<input className={input} value={brand} onChange={(e) => setBrand(e.target.value)} />
</div>
<div>
<label className={label}></label>
<select className={input} value={categoryID} onChange={(e) => setCategoryID(e.target.value)}>
<option value=""></option>
{categories.map((c) => (
<option key={c.id} value={c.id}>
{c.name_zh} ({c.path})
</option>
))}
</select>
</div>
<div>
<label className={label}></label>
<input className={input} value={country} onChange={(e) => setCountry(e.target.value)} />
</div>
<div>
<label className={label}></label>
<input
type="number"
step="any"
className={input}
value={netValue}
onChange={(e) => setNetValue(e.target.value)}
/>
</div>
<div>
<label className={label}> (g/ml/cl)</label>
<input className={input} value={netUnit} onChange={(e) => setNetUnit(e.target.value)} />
</div>
</div>
</div>
<div className="bg-white border rounded-lg p-5 mt-4">
<h2 className="font-medium text-gray-700 mb-3"></h2>
<label className={label}></label>
<textarea
className={input}
rows={3}
value={ingredients}
onChange={(e) => setIngredients(e.target.value)}
/>
<div className="mt-3">
<label className={label}></label>
<select className={input} value={basis} onChange={(e) => setBasis(e.target.value)}>
<option value=""></option>
<option value="per_100g"> 100g</option>
<option value="per_100ml"> 100ml</option>
<option value="per_serving"></option>
</select>
</div>
<div className="grid grid-cols-2 sm:grid-cols-4 gap-3 mt-3">
{NUTRI_FIELDS.map((f) => (
<div key={f.key}>
<label className={label}>{f.label}</label>
<input
type="number"
step="any"
className={input}
value={nutri[f.key] || ""}
onChange={(e) => setNutri({ ...nutri, [f.key]: e.target.value })}
/>
</div>
))}
</div>
</div>
<div className="bg-white border rounded-lg p-5 mt-4">
<h2 className="font-medium text-gray-700 mb-3"> URL</h2>
{images.length > 0 && (
<ul className="mb-3 space-y-1">
{images.map((im, i) => (
<li key={i} className="flex items-center gap-2 text-sm">
<span className="truncate text-gray-600 flex-1">{im.url}</span>
<button
type="button"
onClick={() => setImages(images.filter((_, j) => j !== i))}
className="text-gray-400 hover:text-red-500"
>
<Trash2 className="w-4 h-4" />
</button>
</li>
))}
</ul>
)}
<div className="flex gap-2">
<input
className={input}
placeholder="图片 URL"
value={imageURL}
onChange={(e) => setImageURL(e.target.value)}
/>
<button
type="button"
onClick={() => {
if (imageURL.trim()) {
setImages([...images, { url: imageURL.trim(), kind: "front" }]);
setImageURL("");
}
}}
className="shrink-0 px-3 rounded-md border text-sm flex items-center gap-1 hover:bg-gray-50"
>
<PlusCircle className="w-4 h-4" />
</button>
</div>
</div>
<div className="bg-white border rounded-lg p-5 mt-4">
<h2 className="font-medium text-gray-700 mb-3"></h2>
<div className="grid grid-cols-1 sm:grid-cols-2 gap-4">
<div>
<label className={label}></label>
<input className={input} value={submitter} onChange={(e) => setSubmitter(e.target.value)} />
</div>
<div>
<label className={label}>/</label>
<input className={input} value={contact} onChange={(e) => setContact(e.target.value)} />
</div>
<div className="sm:col-span-2">
<label className={label}> / </label>
<input className={input} value={note} onChange={(e) => setNote(e.target.value)} />
</div>
</div>
</div>
<div className="mt-5 flex items-center gap-3">
<button
type="submit"
disabled={submitting}
className="px-6 py-2.5 rounded-lg bg-emerald-600 text-white font-medium hover:bg-emerald-700 disabled:opacity-60"
>
{submitting ? "提交中…" : "提交审核"}
</button>
<button type="button" onClick={onDone} className="text-sm text-gray-500 hover:underline">
</button>
</div>
</form>
);
}
+118
View File
@@ -0,0 +1,118 @@
import { useState } from "react";
import { Search, PlusCircle, Code2 } from "lucide-react";
import { api } from "../api";
import type { ProductSummary } from "../types";
export default function Home({
onOpen,
onContribute,
onApi,
}: {
onOpen: (id: string) => void;
onContribute: () => void;
onApi: () => void;
}) {
const [q, setQ] = useState("");
const [items, setItems] = useState<ProductSummary[]>([]);
const [total, setTotal] = useState(0);
const [searched, setSearched] = useState(false);
const [loading, setLoading] = useState(false);
const [error, setError] = useState("");
async function run(e?: React.FormEvent) {
e?.preventDefault();
setLoading(true);
setError("");
try {
const res = await api.search(q.trim(), 1, 30);
setItems(res.items);
setTotal(res.total);
setSearched(true);
} catch (err) {
setError(err instanceof Error ? err.message : "搜索失败");
} finally {
setLoading(false);
}
}
return (
<div>
<div className="text-center py-10">
<h1 className="text-3xl font-bold text-gray-800"></h1>
<p className="mt-2 text-gray-500">
</p>
<form onSubmit={run} className="mt-6 max-w-2xl mx-auto flex gap-2">
<div className="flex-1 flex items-center gap-2 bg-white border rounded-lg px-3 shadow-sm focus-within:ring-2 focus-within:ring-emerald-400">
<Search className="w-5 h-5 text-gray-400" />
<input
autoFocus
value={q}
onChange={(e) => setQ(e.target.value)}
placeholder="例如:可乐、Nutella、5449000000996"
className="flex-1 py-3 outline-none bg-transparent"
/>
</div>
<button
type="submit"
disabled={loading}
className="px-6 rounded-lg bg-emerald-600 text-white font-medium hover:bg-emerald-700 disabled:opacity-60"
>
{loading ? "检索中…" : "检索"}
</button>
</form>
<button
onClick={onApi}
className="mt-4 inline-flex items-center gap-1.5 text-sm text-emerald-700 hover:underline"
>
<Code2 className="w-4 h-4" /> API
</button>
</div>
{error && (
<div className="max-w-2xl mx-auto bg-red-50 text-red-700 text-sm rounded-md px-4 py-2">
{error}
</div>
)}
{searched && (
<div className="mt-2">
<div className="text-sm text-gray-500 mb-2">
{total} {q ? `(关键词:${q}` : ""}
</div>
{items.length === 0 ? (
<div className="bg-white border rounded-lg p-8 text-center text-gray-500">
<p></p>
<button
onClick={onContribute}
className="mt-3 inline-flex items-center gap-1.5 text-emerald-700 hover:underline"
>
<PlusCircle className="w-4 h-4" />
</button>
</div>
) : (
<ul className="bg-white border rounded-lg divide-y">
{items.map((p) => (
<li key={p.id}>
<button
onClick={() => onOpen(p.id)}
className="w-full text-left px-4 py-3 hover:bg-gray-50 flex items-center justify-between gap-4"
>
<div>
<div className="font-medium text-gray-800">{p.name}</div>
<div className="text-xs text-gray-500 mt-0.5">
{p.brand || "未知品牌"}
{p.gtin ? ` · ${p.gtin}` : ""}
</div>
</div>
<span className="text-xs text-gray-400">{p.category_path || ""}</span>
</button>
</li>
))}
</ul>
)}
</div>
)}
</div>
);
}
@@ -0,0 +1,104 @@
import { useEffect, useState } from "react";
import { ArrowLeft } from "lucide-react";
import { api } from "../api";
import type { Product } from "../types";
import { NUTRIMENT_LABELS } from "../types";
function Row({ label, value }: { label: string; value: React.ReactNode }) {
if (value === null || value === undefined || value === "") return null;
return (
<div className="flex py-2 border-b last:border-0 text-sm">
<div className="w-32 shrink-0 text-gray-400">{label}</div>
<div className="text-gray-800">{value}</div>
</div>
);
}
export default function ProductView({ id, onBack }: { id: string; onBack: () => void }) {
const [p, setP] = useState<Product | null>(null);
const [error, setError] = useState("");
useEffect(() => {
api.product(id).then(setP).catch((e) => setError(e.message));
}, [id]);
if (error) {
return (
<div>
<button onClick={onBack} className="text-sm text-gray-500 flex items-center gap-1 mb-4">
<ArrowLeft className="w-4 h-4" />
</button>
<div className="bg-red-50 text-red-700 text-sm rounded-md px-4 py-3">{error}</div>
</div>
);
}
if (!p) return <div className="text-gray-400"></div>;
const basisLabel: Record<string, string> = {
per_100g: "每 100g",
per_100ml: "每 100ml",
per_serving: "每份",
};
const nutriEntries = Object.entries(p.nutriments || {}).filter(
([, v]) => v !== null && v !== undefined,
);
return (
<div>
<button onClick={onBack} className="text-sm text-gray-500 flex items-center gap-1 mb-4">
<ArrowLeft className="w-4 h-4" />
</button>
<div className="bg-white border rounded-lg p-5">
<div className="flex items-start justify-between gap-4">
<h1 className="text-xl font-semibold text-gray-800">{p.name}</h1>
<span className="shrink-0 text-xs bg-emerald-50 text-emerald-700 rounded px-2 py-1">
{Math.round(p.quality_score * 100)}
</span>
</div>
<div className="mt-4">
<Row label="品牌" value={p.brand} />
<Row label="条码 (GTIN)" value={p.gtin} />
<Row label="品类" value={p.category_path} />
<Row
label="净含量"
value={
p.net_content_value != null
? `${p.net_content_value} ${p.net_content_unit || ""}`
: null
}
/>
<Row label="产地" value={p.country_of_origin} />
<Row label="Nutri-Score" value={p.nutri_score} />
<Row label="配料" value={p.ingredients_text} />
<Row
label="过敏原"
value={p.allergens && p.allergens.length ? p.allergens.join("、") : null}
/>
<Row
label="添加剂"
value={p.additives && p.additives.length ? p.additives.join("、") : null}
/>
</div>
</div>
{nutriEntries.length > 0 && (
<div className="bg-white border rounded-lg p-5 mt-4">
<h2 className="font-medium text-gray-700 mb-2">
{p.nutrition_basis ? `${basisLabel[p.nutrition_basis] || p.nutrition_basis}` : ""}
</h2>
<div className="grid grid-cols-2 sm:grid-cols-3 gap-2 text-sm">
{nutriEntries.map(([k, v]) => (
<div key={k} className="bg-gray-50 rounded px-3 py-2">
<div className="text-gray-400 text-xs">{NUTRIMENT_LABELS[k] || k}</div>
<div className="text-gray-800">{String(v)}</div>
</div>
))}
</div>
</div>
)}
</div>
);
}
+16
View File
@@ -0,0 +1,16 @@
@tailwind base;
@tailwind components;
@tailwind utilities;
html,
body,
#root {
height: 100%;
}
body {
margin: 0;
background: #f3f4f6;
font-family: system-ui, -apple-system, "Segoe UI", Roboto, "Helvetica Neue",
Arial, "PingFang SC", "Microsoft YaHei", sans-serif;
}
+10
View File
@@ -0,0 +1,10 @@
import React from "react";
import ReactDOM from "react-dom/client";
import App from "./App";
import "./index.css";
ReactDOM.createRoot(document.getElementById("root")!).render(
<React.StrictMode>
<App />
</React.StrictMode>,
);
+70
View File
@@ -0,0 +1,70 @@
export interface ProductSummary {
id: string;
gtin: string | null;
name: string;
brand: string | null;
category_path: string | null;
}
export interface Product {
id: string;
gtin: string | null;
name: string;
brand: string | null;
category_path: string | null;
net_content_value: number | null;
net_content_unit: string | null;
country_of_origin: string | null;
quality_score: number;
nutriments?: Record<string, unknown> | null;
nutrition_basis?: string | null;
nutri_score?: string | null;
ingredients_text?: string | null;
allergens?: string[] | null;
additives?: string[] | null;
}
export interface Category {
id: string;
name_zh: string;
name_en: string | null;
path: string;
level: number;
}
export interface SubmissionImage {
url: string;
kind: string;
}
export interface SubmissionInput {
gtin?: string | null;
name: string;
brand_name?: string | null;
category_id?: string | null;
net_content_value?: number | null;
net_content_unit?: string | null;
country_of_origin?: string | null;
ingredients_text?: string | null;
nutriments?: Record<string, number> | null;
nutrition_basis?: string | null;
serving_size?: string | null;
nutri_score?: string | null;
images?: SubmissionImage[];
submitter_name?: string | null;
submitter_contact?: string | null;
note?: string | null;
}
export const NUTRIMENT_LABELS: Record<string, string> = {
energy_kcal: "能量 (kcal)",
energy_kj: "能量 (kJ)",
fat: "脂肪 (g)",
saturated_fat: "饱和脂肪 (g)",
carbohydrates: "碳水 (g)",
sugars: "糖 (g)",
proteins: "蛋白质 (g)",
salt: "盐 (g)",
sodium: "钠 (g)",
fiber: "膳食纤维 (g)",
};
+1
View File
@@ -0,0 +1 @@
/// <reference types="vite/client" />
+6
View File
@@ -0,0 +1,6 @@
/** @type {import('tailwindcss').Config} */
export default {
content: ["./index.html", "./src/**/*.{ts,tsx}"],
theme: { extend: {} },
plugins: [],
};
+21
View File
@@ -0,0 +1,21 @@
{
"compilerOptions": {
"target": "ES2020",
"useDefineForClassFields": true,
"lib": ["ES2020", "DOM", "DOM.Iterable"],
"module": "ESNext",
"skipLibCheck": true,
"moduleResolution": "bundler",
"allowImportingTsExtensions": true,
"resolveJsonModule": true,
"isolatedModules": true,
"noEmit": true,
"jsx": "react-jsx",
"strict": true,
"noUnusedLocals": true,
"noUnusedParameters": true,
"noFallthroughCasesInSwitch": true
},
"include": ["src"],
"references": [{ "path": "./tsconfig.node.json" }]
}
+11
View File
@@ -0,0 +1,11 @@
{
"compilerOptions": {
"composite": true,
"skipLibCheck": true,
"module": "ESNext",
"moduleResolution": "bundler",
"allowSyntheticDefaultImports": true,
"strict": true
},
"include": ["vite.config.ts"]
}
+9
View File
@@ -0,0 +1,9 @@
import { defineConfig } from "vite";
import react from "@vitejs/plugin-react";
// Served at the site root by the public read-only Go binary.
export default defineConfig({
base: "/",
plugins: [react()],
build: { outDir: "dist", emptyOutDir: true },
});