Compare commits

..

6 Commits

Author SHA1 Message Date
novaalphastrikeomegaz663 2820823b36 feat(api): API keys + Redis rate limiting + usage stats
CI / Python (ingestion) (pull_request) Successful in 12s
CI / Migrations (postgres) (pull_request) Successful in 24s
CI / Go (api) (pull_request) Successful in 53s
Add an optional API-key layer to the public read-only API. Keys grant
higher per-minute rate limits and attribute usage; anonymous callers are
still allowed at a lower IP-based budget.

- migration 0008_api_key: api_key table (sha256 hash only, plaintext shown once)
- apikey pkg: key generation + hashing
- ratelimit pkg: Redis fixed-window limiter + per-key usage counters; fails open
- public API middleware: X-API-Key / Bearer auth, X-RateLimit-* headers, 429+Retry-After
- admin: issue/list/revoke keys + usage view (API + UI tab)

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 08:24:07 +00:00
lixu 7d7a8f1baf Merge pull request 'feat(ingestion): harden OFF ingestion for bulk seeding' (#6) from devin/1781938549-ingestion-resilience into main
CI / Python (ingestion) (push) Successful in 11s
CI / Migrations (postgres) (push) Successful in 22s
CI / Go (api) (push) Successful in 36s
2026-06-20 15:54:07 +08:00
lixu e34e007f28 Merge pull request 'feat(barcode): 多条码管理(迁移 + GTIN 校验 + 后台/公开 API + 后台 UI)' (#5) from devin/1781934649-multi-barcode into main
CI / Python (ingestion) (push) Successful in 11s
CI / Migrations (postgres) (push) Successful in 22s
CI / Go (api) (push) Successful in 37s
2026-06-20 15:53:43 +08:00
novaalphastrikeomegaz663 c090bcdc9b style(ingest): ruff format country-seed tests
CI / Python (ingestion) (pull_request) Successful in 9s
CI / Migrations (postgres) (pull_request) Successful in 13s
CI / Go (api) (pull_request) Successful in 29s
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 07:32:30 +00:00
novaalphastrikeomegaz663 2255243081 feat(ingest): country-focused seeding (collect domestic CN products)
CI / Python (ingestion) (pull_request) Failing after 6s
CI / Migrations (postgres) (pull_request) Successful in 16s
CI / Go (api) (pull_request) Failing after 11m35s
Add a market-focused seeding path so the catalogue can be built from
domestic products rather than the English-heavy global default:

- adapter.fetch_by_country(country): OFF search filtered by
  countries_tags_en, sorted by unique_scans_n (most-scanned first),
  de-duplicated across pages since OFF popularity ordering is unstable.
- is_cn_gs1(code): True for GS1-China company prefixes (690-699),
  i.e. genuinely domestic items vs. imports merely sold in China.
- seed_off --country <slug> [--domestic-only] [--page-size/--max-pages]:
  e.g. 'seed_off --country china --domestic-only' loads only 69x
  barcodes.

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 07:14:45 +00:00
novaalphastrikeomegaz663 044c870df7 feat(ingestion): harden OFF ingestion for bulk seeding
CI / Go (api) (pull_request) Failing after 22s
CI / Python (ingestion) (pull_request) Successful in 14s
CI / Migrations (postgres) (pull_request) Failing after 18s
- Add retry/backoff (429 + 5xx, Retry-After aware) to the OFF adapter so
  transient API errors no longer abort a run.
- Clamp bounded text fields (serving_size, net_content_unit,
  country_of_origin) to their column widths in transform; long OFF values
  previously raised StringDataRightTruncation and rolled back the batch.
- Load each record inside a savepoint (load_record_safe) so one malformed
  source record is skipped instead of aborting the whole import; jobs now
  report an errored count.
- Tests for retry behaviour, serving_size clamping, and per-record isolation.

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-06-20 06:55:56 +00:00
31 changed files with 1539 additions and 39 deletions
+16 -3
View File
@@ -4,9 +4,10 @@ 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";
import ApiKeysPage from "./components/ApiKeysPage";
import { Inbox, KeyRound, LogOut, Package } from "lucide-react";
type Tab = "products" | "submissions";
type Tab = "products" | "submissions" | "keys";
type View = { name: "list" } | { name: "detail"; id: string };
export default function App() {
@@ -99,6 +100,16 @@ export default function App() {
</span>
)}
</button>
<button
onClick={() => setTab("keys")}
className={`px-3 py-1.5 rounded-md flex items-center gap-1.5 ${
tab === "keys"
? "bg-emerald-50 text-emerald-700"
: "text-gray-600 hover:bg-gray-100"
}`}
>
<KeyRound className="h-4 w-4" /> API
</button>
</nav>
</div>
<div className="flex items-center gap-4 text-sm text-gray-600">
@@ -112,7 +123,9 @@ export default function App() {
</div>
</header>
<main className="flex-1 overflow-auto p-6">
{tab === "submissions" ? (
{tab === "keys" ? (
<ApiKeysPage />
) : tab === "submissions" ? (
<SubmissionsPage onPending={setPending} />
) : view.name === "list" ? (
<ProductList onOpen={(id) => setView({ name: "detail", id })} />
+14
View File
@@ -123,4 +123,18 @@ export const api = {
method: "POST",
body: JSON.stringify({ note }),
}),
listApiKeys: () =>
request<{ items: import("./types").ApiKey[] }>("/keys"),
createApiKey: (body: {
name: string;
owner_email?: string;
tier?: string;
rate_limit_per_min?: number;
}) =>
request<{ key: string; item: import("./types").ApiKey; warning: string }>(
"/keys",
{ method: "POST", body: JSON.stringify(body) },
),
revokeApiKey: (id: string) =>
request<{ status: string }>(`/keys/${id}`, { method: "DELETE" }),
};
@@ -0,0 +1,278 @@
import { useEffect, useState } from "react";
import { api, ApiError } from "../api";
import type { ApiKey } from "../types";
import { Copy, KeyRound, Plus, Trash2 } from "lucide-react";
const TIERS = [
{ key: "free", label: "免费 (free)", rate: 120 },
{ key: "partner", label: "合作方 (partner)", rate: 600 },
{ key: "internal", label: "内部 (internal)", rate: 6000 },
];
function tierLabel(tier: string): string {
return TIERS.find((t) => t.key === tier)?.label ?? tier;
}
export default function ApiKeysPage() {
const [rows, setRows] = useState<ApiKey[]>([]);
const [error, setError] = useState("");
const [creating, setCreating] = useState(false);
const [newKey, setNewKey] = useState<string | null>(null);
async function load() {
setError("");
try {
const res = await api.listApiKeys();
setRows(res.items);
} catch (e) {
setError(e instanceof ApiError ? e.message : "加载失败");
}
}
useEffect(() => {
load();
}, []);
async function revoke(id: string, name: string) {
if (!confirm(`确认吊销密钥「${name}」?使用该密钥的请求将立即被拒绝。`)) return;
try {
await api.revokeApiKey(id);
await load();
} catch (e) {
setError(e instanceof ApiError ? e.message : "操作失败");
}
}
return (
<div className="max-w-4xl">
<div className="flex items-center justify-between mb-4">
<div>
<h2 className="text-lg font-semibold text-gray-800 flex items-center gap-2">
<KeyRound className="h-5 w-5 text-emerald-600" /> API
</h2>
<p className="text-sm text-gray-500 mt-1">
API
</p>
</div>
<button
onClick={() => setCreating(true)}
className="px-4 py-2 rounded-lg bg-emerald-600 text-white text-sm font-medium hover:bg-emerald-700 flex items-center gap-1.5"
>
<Plus className="h-4 w-4" />
</button>
</div>
{error && (
<div className="mb-3 bg-red-50 text-red-700 text-sm rounded px-4 py-2">{error}</div>
)}
{newKey && (
<div className="mb-4 bg-amber-50 border border-amber-200 rounded-lg p-4">
<div className="text-sm font-medium text-amber-800 mb-1">
</div>
<div className="flex items-center gap-2">
<code className="flex-1 bg-white border rounded px-3 py-2 text-sm break-all">
{newKey}
</code>
<button
onClick={() => navigator.clipboard?.writeText(newKey)}
className="px-3 py-2 rounded border text-sm text-gray-600 hover:bg-gray-50 flex items-center gap-1"
>
<Copy className="h-4 w-4" />
</button>
<button
onClick={() => setNewKey(null)}
className="px-3 py-2 rounded text-sm text-gray-500 hover:bg-gray-100"
>
</button>
</div>
</div>
)}
{creating && (
<CreateKeyForm
onClose={() => setCreating(false)}
onCreated={(plaintext) => {
setCreating(false);
setNewKey(plaintext);
load();
}}
/>
)}
<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>
<th className="px-4 py-2 font-medium"></th>
</tr>
</thead>
<tbody className="divide-y">
{rows.length === 0 ? (
<tr>
<td colSpan={7} className="px-4 py-8 text-center text-gray-400">
</td>
</tr>
) : (
rows.map((k) => (
<tr key={k.id} className={k.revoked_at ? "opacity-50" : ""}>
<td className="px-4 py-2 text-gray-800">
{k.name}
{k.owner_email && (
<span className="block text-xs text-gray-400">{k.owner_email}</span>
)}
</td>
<td className="px-4 py-2 text-gray-500">
<code>{k.key_prefix}</code>
</td>
<td className="px-4 py-2 text-gray-600">{tierLabel(k.tier)}</td>
<td className="px-4 py-2 text-gray-600">{k.rate_limit_per_min}</td>
<td className="px-4 py-2 text-gray-600">
{k.usage.today} / {k.usage.total}
</td>
<td className="px-4 py-2">
{k.revoked_at ? (
<span className="text-xs rounded px-2 py-0.5 bg-red-50 text-red-700">
</span>
) : (
<span className="text-xs rounded px-2 py-0.5 bg-emerald-50 text-emerald-700">
</span>
)}
</td>
<td className="px-4 py-2 text-right">
{!k.revoked_at && (
<button
onClick={() => revoke(k.id, k.name)}
className="text-gray-400 hover:text-red-600"
title="吊销"
>
<Trash2 className="h-4 w-4" />
</button>
)}
</td>
</tr>
))
)}
</tbody>
</table>
</div>
</div>
);
}
function CreateKeyForm({
onClose,
onCreated,
}: {
onClose: () => void;
onCreated: (plaintext: string) => void;
}) {
const [name, setName] = useState("");
const [ownerEmail, setOwnerEmail] = useState("");
const [tier, setTier] = useState("free");
const [rate, setRate] = useState(120);
const [busy, setBusy] = useState(false);
const [error, setError] = useState("");
function pickTier(t: string) {
setTier(t);
const def = TIERS.find((x) => x.key === t);
if (def) setRate(def.rate);
}
async function submit() {
if (!name.trim()) {
setError("名称不能为空");
return;
}
setBusy(true);
setError("");
try {
const res = await api.createApiKey({
name: name.trim(),
owner_email: ownerEmail.trim() || undefined,
tier,
rate_limit_per_min: rate,
});
onCreated(res.key);
} catch (e) {
setError(e instanceof ApiError ? e.message : "创建失败");
} finally {
setBusy(false);
}
}
return (
<div className="mb-4 bg-white border rounded-lg p-5">
<h3 className="font-medium text-gray-700 mb-3"></h3>
{error && <div className="mb-3 bg-red-50 text-red-700 text-sm rounded px-3 py-2">{error}</div>}
<div className="grid grid-cols-2 gap-4">
<label className="block">
<span className="text-xs text-gray-500"> *</span>
<input
className="w-full border rounded-md px-3 py-2 text-sm mt-1"
value={name}
onChange={(e) => setName(e.target.value)}
placeholder="例如:我的 App / 合作方 X"
/>
</label>
<label className="block">
<span className="text-xs text-gray-500"></span>
<input
className="w-full border rounded-md px-3 py-2 text-sm mt-1"
value={ownerEmail}
onChange={(e) => setOwnerEmail(e.target.value)}
placeholder="owner@example.com"
/>
</label>
<label className="block">
<span className="text-xs text-gray-500"></span>
<select
className="w-full border rounded-md px-3 py-2 text-sm mt-1 bg-white"
value={tier}
onChange={(e) => pickTier(e.target.value)}
>
{TIERS.map((t) => (
<option key={t.key} value={t.key}>
{t.label}
</option>
))}
</select>
</label>
<label className="block">
<span className="text-xs text-gray-500">/</span>
<input
type="number"
min={1}
className="w-full border rounded-md px-3 py-2 text-sm mt-1"
value={rate}
onChange={(e) => setRate(Math.max(1, parseInt(e.target.value || "1", 10)))}
/>
</label>
</div>
<div className="mt-4 flex gap-2">
<button
onClick={submit}
disabled={busy}
className="px-4 py-2 rounded bg-emerald-600 text-white text-sm hover:bg-emerald-700 disabled:opacity-60"
>
</button>
<button onClick={onClose} className="px-4 py-2 rounded border text-sm text-gray-600">
</button>
</div>
</div>
);
}
+19
View File
@@ -137,6 +137,25 @@ export interface SubmissionDetail {
existing_product?: ProductDetail;
}
export interface ApiKeyUsage {
total: number;
today: number;
last_used_at?: number | null;
}
export interface ApiKey {
id: string;
name: string;
key_prefix: string;
owner_email: string | null;
tier: string;
rate_limit_per_min: number;
revoked_at: string | null;
created_by: string | null;
created_at: string;
usage: ApiKeyUsage;
}
export const FIELD_LABELS: Record<string, string> = {
name: "名称",
gtin: "条码",
+4 -1
View File
@@ -17,6 +17,7 @@ import (
"github.com/baicai2026-baicai/goods/api/internal/adminstore"
"github.com/baicai2026-baicai/goods/api/internal/adminweb"
"github.com/baicai2026-baicai/goods/api/internal/auth"
"github.com/baicai2026-baicai/goods/api/internal/ratelimit"
)
func getenv(key, fallback string) string {
@@ -69,7 +70,9 @@ func main() {
}
authn := auth.New(username, passwordHash, secret, 12*time.Hour)
h := adminhandler.New(adminstore.New(pool), authn, basePath, adminweb.Dist())
usage := ratelimit.New(getenv("OPENGOODS_REDIS_URL", "redis://localhost:6379/0"))
h := adminhandler.New(adminstore.New(pool), authn, basePath, adminweb.Dist()).
WithUsage(usage)
srv := &http.Server{
Addr: addr,
+7 -1
View File
@@ -12,6 +12,7 @@ import (
"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/ratelimit"
"github.com/baicai2026-baicai/goods/api/internal/store"
)
@@ -31,7 +32,12 @@ func main() {
log.Printf("warning: database not reachable at startup: %v", err)
}
h := handler.New(store.New(pool), publicweb.Dist())
limiter := ratelimit.New(cfg.RedisURL)
if !limiter.Enabled() {
log.Print("warning: Redis not configured; public API rate limiting disabled")
}
h := handler.New(store.New(pool), publicweb.Dist()).
WithRateLimit(limiter, cfg.AnonRateLimitPerMin)
srv := &http.Server{
Addr: cfg.Addr,
+4
View File
@@ -5,13 +5,17 @@ go 1.23.4
require (
github.com/go-chi/chi/v5 v5.1.0
github.com/jackc/pgx/v5 v5.7.2
github.com/redis/go-redis/v9 v9.18.0
golang.org/x/crypto v0.31.0
)
require (
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect
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
go.uber.org/atomic v1.11.0 // indirect
golang.org/x/sync v0.10.0 // indirect
golang.org/x/text v0.21.0 // indirect
)
+16
View File
@@ -1,6 +1,14 @@
github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs=
github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c=
github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA=
github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0=
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
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/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f h1:lO4WD4F/rVNCu3HqELle0jiPLLBs70cWOduZpkS1E78=
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc=
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=
@@ -11,13 +19,21 @@ 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/klauspost/cpuid/v2 v2.0.9 h1:lgaqFMSdTdQYdZ04uHyN2d/eKdOMyi2YLSvlQIBFYa4=
github.com/klauspost/cpuid/v2 v2.0.9/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg=
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/redis/go-redis/v9 v9.18.0 h1:pMkxYPkEbMPwRdenAzUNyFNrDgHx9U+DrBabWNfSRQs=
github.com/redis/go-redis/v9 v9.18.0/go.mod h1:k3ufPphLU5YXwNTUcCRXGxUoF1fqxnhFQmscfkCoDA0=
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=
github.com/zeebo/xxh3 v1.0.2 h1:xZmwmqxHZA8AI603jOQ0tMqmBr9lPeFwGg6d+xy9DC0=
github.com/zeebo/xxh3 v1.0.2/go.mod h1:5NWz9Sef7zIDm2JHfFlcQvNekmcEl9ekUZQQKCYaDcA=
go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE=
go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0=
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=
+70
View File
@@ -0,0 +1,70 @@
package adminhandler
import (
"encoding/json"
"net/http"
"strings"
"github.com/go-chi/chi/v5"
"github.com/baicai2026-baicai/goods/api/internal/adminstore"
"github.com/baicai2026-baicai/goods/api/internal/auth"
"github.com/baicai2026-baicai/goods/api/internal/ratelimit"
)
// apiKeyView is an issued key plus its usage counters.
type apiKeyView struct {
adminstore.APIKeyRow
Usage ratelimit.UsageStat `json:"usage"`
}
// ListAPIKeys returns all issued keys with usage stats merged in.
func (h *Handler) ListAPIKeys(w http.ResponseWriter, r *http.Request) {
keys, err := h.store.ListAPIKeys(r.Context())
if h.handleErr(w, err) {
return
}
views := make([]apiKeyView, 0, len(keys))
for _, k := range keys {
v := apiKeyView{APIKeyRow: k}
if h.usage != nil {
v.Usage = h.usage.Usage(r.Context(), k.ID)
}
views = append(views, v)
}
writeJSON(w, http.StatusOK, map[string]any{"items": views})
}
// CreateAPIKey issues a new key and returns its plaintext exactly once.
func (h *Handler) CreateAPIKey(w http.ResponseWriter, r *http.Request) {
var in adminstore.APIKeyInput
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
}
if in.Tier != "" && in.Tier != "free" && in.Tier != "partner" && in.Tier != "internal" {
writeError(w, http.StatusBadRequest, "bad_request", "tier 取值无效")
return
}
plaintext, row, err := h.store.CreateAPIKey(r.Context(), in, auth.UserFrom(r.Context()))
if h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusCreated, map[string]any{
"key": plaintext,
"item": row,
"warning": "请立即复制保存此密钥,它只显示这一次,无法再次查看。",
})
}
// RevokeAPIKey disables a key. Subsequent requests with it are rejected.
func (h *Handler) RevokeAPIKey(w http.ResponseWriter, r *http.Request) {
if err := h.store.RevokeAPIKey(r.Context(), chi.URLParam(r, "id")); h.handleErr(w, err) {
return
}
writeJSON(w, http.StatusOK, map[string]string{"status": "revoked"})
}
+13
View File
@@ -16,6 +16,7 @@ import (
"github.com/baicai2026-baicai/goods/api/internal/adminstore"
"github.com/baicai2026-baicai/goods/api/internal/auth"
"github.com/baicai2026-baicai/goods/api/internal/gtin"
"github.com/baicai2026-baicai/goods/api/internal/ratelimit"
)
// Handler holds the admin dependencies.
@@ -25,6 +26,7 @@ type Handler struct {
basePath string
spa fs.FS
submitLimit *rateLimiter
usage *ratelimit.Limiter
}
// New constructs an admin Handler. basePath is e.g. "/ping" (no trailing slash).
@@ -39,6 +41,13 @@ func New(store *adminstore.Store, authn *auth.Authenticator, basePath string, sp
}
}
// WithUsage attaches a Redis-backed limiter used to read per-key usage counters
// for the API-key management view. Optional; without it usage shows as zero.
func (h *Handler) WithUsage(l *ratelimit.Limiter) *Handler {
h.usage = l
return h
}
// Router builds the HTTP handler.
func (h *Handler) Router() http.Handler {
r := chi.NewRouter()
@@ -73,6 +82,10 @@ func (h *Handler) Router() http.Handler {
r.Get("/api/submissions/{id}", h.GetSubmission)
r.Post("/api/submissions/{id}/approve", h.ApproveSubmission)
r.Post("/api/submissions/{id}/reject", h.RejectSubmission)
r.Get("/api/keys", h.ListAPIKeys)
r.Post("/api/keys", h.CreateAPIKey)
r.Delete("/api/keys/{id}", h.RevokeAPIKey)
})
r.Handle("/*", http.HandlerFunc(h.serveSPA))
+130
View File
@@ -0,0 +1,130 @@
package adminstore
import (
"context"
"errors"
"strings"
"time"
"github.com/jackc/pgx/v5/pgconn"
"github.com/baicai2026-baicai/goods/api/internal/apikey"
)
// APIKeyRow is an admin-facing view of an issued API key (never the secret).
type APIKeyRow struct {
ID string `json:"id"`
Name string `json:"name"`
KeyPrefix string `json:"key_prefix"`
OwnerEmail *string `json:"owner_email"`
Tier string `json:"tier"`
RateLimitPerMin int `json:"rate_limit_per_min"`
RevokedAt *string `json:"revoked_at"`
CreatedBy *string `json:"created_by"`
CreatedAt string `json:"created_at"`
}
// APIKeyInput holds the fields accepted when issuing a key.
type APIKeyInput struct {
Name string `json:"name"`
OwnerEmail string `json:"owner_email"`
Tier string `json:"tier"`
RateLimitPerMin int `json:"rate_limit_per_min"`
}
// CreateAPIKey issues a new key, returning the one-time plaintext alongside the
// stored row. Only the SHA-256 hash and a short display prefix are persisted.
func (s *Store) CreateAPIKey(ctx context.Context, in APIKeyInput, createdBy string) (plaintext string, row APIKeyRow, err error) {
tier := in.Tier
if tier == "" {
tier = "free"
}
rate := in.RateLimitPerMin
if rate <= 0 {
rate = 120
}
var owner *string
if e := strings.TrimSpace(in.OwnerEmail); e != "" {
owner = &e
}
key, hash, prefix, err := apikey.Generate()
if err != nil {
return "", row, err
}
var revoked, created *time.Time
var createdByOut *string
err = s.pool.QueryRow(ctx, `
INSERT INTO api_key (name, key_prefix, key_hash, owner_email, tier, rate_limit_per_min, created_by)
VALUES ($1, $2, $3, $4, $5, $6, $7)
RETURNING id, name, key_prefix, owner_email, tier, rate_limit_per_min, revoked_at, created_by, created_at`,
strings.TrimSpace(in.Name), prefix, hash, owner, tier, rate, createdBy,
).Scan(&row.ID, &row.Name, &row.KeyPrefix, &row.OwnerEmail, &row.Tier,
&row.RateLimitPerMin, &revoked, &createdByOut, &created)
if err != nil {
return "", row, err
}
row.CreatedBy = createdByOut
if created != nil {
row.CreatedAt = created.Format(time.RFC3339)
}
return key, row, nil
}
// ListAPIKeys returns all keys (active first, newest first).
func (s *Store) ListAPIKeys(ctx context.Context) ([]APIKeyRow, error) {
rows, err := s.pool.Query(ctx, `
SELECT id, name, key_prefix, owner_email, tier, rate_limit_per_min, revoked_at, created_by, created_at
FROM api_key
ORDER BY (revoked_at IS NULL) DESC, created_at DESC`)
if err != nil {
return nil, err
}
defer rows.Close()
out := []APIKeyRow{}
for rows.Next() {
var r APIKeyRow
var revoked, created *time.Time
if err := rows.Scan(&r.ID, &r.Name, &r.KeyPrefix, &r.OwnerEmail, &r.Tier,
&r.RateLimitPerMin, &revoked, &r.CreatedBy, &created); err != nil {
return nil, err
}
if revoked != nil {
v := revoked.Format(time.RFC3339)
r.RevokedAt = &v
}
if created != nil {
r.CreatedAt = created.Format(time.RFC3339)
}
out = append(out, r)
}
return out, rows.Err()
}
// RevokeAPIKey marks a key revoked. Revoking an already-revoked or missing key
// returns ErrNotFound.
func (s *Store) RevokeAPIKey(ctx context.Context, id string) error {
tag, err := s.pool.Exec(ctx,
"UPDATE api_key SET revoked_at = now() WHERE id = $1 AND revoked_at IS NULL", id)
if err != nil {
if isInvalidUUID(err) {
return ErrNotFound
}
return err
}
if tag.RowsAffected() == 0 {
return ErrNotFound
}
return nil
}
// isInvalidUUID reports whether err is a Postgres invalid-UUID-text error,
// which happens when a non-UUID id is supplied.
func isInvalidUUID(err error) bool {
var pgErr *pgconn.PgError
if errors.As(err, &pgErr) {
return pgErr.Code == "22P02"
}
return false
}
+49
View File
@@ -0,0 +1,49 @@
// Package apikey handles generation and hashing of public-API keys.
//
// A key looks like "og_live_<random>". Only the SHA-256 hash is ever persisted;
// the plaintext is returned once at creation time and cannot be recovered.
package apikey
import (
"crypto/rand"
"crypto/sha256"
"encoding/base64"
"encoding/hex"
"strings"
)
// Prefix is the human-readable scheme prefix every key carries.
const Prefix = "og_live_"
// prefixLen is how many leading characters (including Prefix) are stored in
// api_key.key_prefix for identifying a key without revealing its secret.
const prefixLen = 12
// Generate returns a new random key (plaintext), its SHA-256 hash, and a short
// display prefix. The plaintext must be shown to the caller exactly once.
func Generate() (key, hash, displayPrefix string, err error) {
buf := make([]byte, 24)
if _, err = rand.Read(buf); err != nil {
return "", "", "", err
}
// URL-safe, no padding => stable, copy-pasteable token body.
body := base64.RawURLEncoding.EncodeToString(buf)
key = Prefix + body
hash = Hash(key)
displayPrefix = key
if len(displayPrefix) > prefixLen {
displayPrefix = displayPrefix[:prefixLen]
}
return key, hash, displayPrefix, nil
}
// Hash returns the hex-encoded SHA-256 of a key, used for storage and lookup.
func Hash(key string) string {
sum := sha256.Sum256([]byte(strings.TrimSpace(key)))
return hex.EncodeToString(sum[:])
}
// Looks like a key issued by this service (cheap pre-check before hashing).
func IsWellFormed(key string) bool {
return strings.HasPrefix(key, Prefix) && len(key) > len(Prefix)+8
}
+60
View File
@@ -0,0 +1,60 @@
package apikey
import "testing"
func TestGenerate(t *testing.T) {
key, hash, prefix, err := Generate()
if err != nil {
t.Fatalf("Generate: %v", err)
}
if !IsWellFormed(key) {
t.Fatalf("generated key not well-formed: %q", key)
}
if Hash(key) != hash {
t.Fatalf("Hash(key) != returned hash")
}
if len(prefix) != prefixLen || key[:prefixLen] != prefix {
t.Fatalf("prefix %q not a %d-char prefix of key %q", prefix, prefixLen, key)
}
if len(hash) != 64 {
t.Fatalf("hash not hex sha-256: %q", hash)
}
}
func TestGenerateUnique(t *testing.T) {
seen := map[string]bool{}
for i := 0; i < 100; i++ {
k, _, _, err := Generate()
if err != nil {
t.Fatal(err)
}
if seen[k] {
t.Fatalf("duplicate key generated: %q", k)
}
seen[k] = true
}
}
func TestHashStableAndTrimmed(t *testing.T) {
if Hash("og_live_abc") != Hash(" og_live_abc ") {
t.Fatal("Hash should ignore surrounding whitespace")
}
if Hash("a") == Hash("b") {
t.Fatal("distinct inputs must hash differently")
}
}
func TestIsWellFormed(t *testing.T) {
cases := map[string]bool{
"og_live_abcdefghijkl": true, // body longer than 8 chars
"og_live_": false, // empty body
"og_live_abc": false, // body too short
"nope_abcdefghijkl": false, // wrong prefix
"": false,
}
for in, want := range cases {
if got := IsWellFormed(in); got != want {
t.Errorf("IsWellFormed(%q) = %v, want %v", in, got, want)
}
}
}
+18 -6
View File
@@ -2,26 +2,38 @@ package config
import (
"os"
"strconv"
)
// Config holds runtime configuration for the OpenGoods API server.
// Values are read from environment variables with sensible defaults so the
// server can boot in a local Docker Compose setup without extra configuration.
type Config struct {
Addr string
DatabaseURL string
RedisURL string
Addr string
DatabaseURL string
RedisURL string
AnonRateLimitPerMin int
}
// Load reads configuration from the environment.
func Load() Config {
return Config{
Addr: getenv("OPENGOODS_ADDR", ":8080"),
DatabaseURL: getenv("OPENGOODS_DATABASE_URL", "postgres://opengoods:opengoods@localhost:5432/opengoods?sslmode=disable"),
RedisURL: getenv("OPENGOODS_REDIS_URL", "redis://localhost:6379/0"),
Addr: getenv("OPENGOODS_ADDR", ":8080"),
DatabaseURL: getenv("OPENGOODS_DATABASE_URL", "postgres://opengoods:opengoods@localhost:5432/opengoods?sslmode=disable"),
RedisURL: getenv("OPENGOODS_REDIS_URL", "redis://localhost:6379/0"),
AnonRateLimitPerMin: getenvInt("OPENGOODS_ANON_RATE_LIMIT_PER_MIN", 60),
}
}
func getenvInt(key string, fallback int) int {
if v, ok := os.LookupEnv(key); ok && v != "" {
if n, err := strconv.Atoi(v); err == nil && n > 0 {
return n
}
}
return fallback
}
func getenv(key, fallback string) string {
if v, ok := os.LookupEnv(key); ok && v != "" {
return v
+23 -3
View File
@@ -15,6 +15,7 @@ import (
"github.com/go-chi/chi/v5"
"github.com/go-chi/chi/v5/middleware"
"github.com/baicai2026-baicai/goods/api/internal/ratelimit"
"github.com/baicai2026-baicai/goods/api/internal/store"
)
@@ -24,17 +25,35 @@ const APIVersion = "v1"
const (
defaultPageSize = 20
maxPageSize = 100
// defaultAnonLimit is the per-minute request budget for unauthenticated
// callers (identified by client IP) when none is configured.
defaultAnonLimit = 60
)
// Handler holds dependencies shared by the HTTP routes.
type Handler struct {
store *store.Store
spa fs.FS
store *store.Store
spa fs.FS
limiter *ratelimit.Limiter
anonLimit int
}
// New constructs a Handler backed by the given store. spa may be nil (JSON-only).
// Rate limiting is disabled until WithRateLimit is called.
func New(s *store.Store, spa fs.FS) *Handler {
return &Handler{store: s, spa: spa}
return &Handler{store: s, spa: spa, anonLimit: defaultAnonLimit}
}
// WithRateLimit attaches a Redis-backed limiter and the anonymous per-minute
// budget, enabling rate limiting + usage tracking on the public API routes.
// A non-positive anonPerMin keeps the default.
func (h *Handler) WithRateLimit(l *ratelimit.Limiter, anonPerMin int) *Handler {
h.limiter = l
if anonPerMin > 0 {
h.anonLimit = anonPerMin
}
return h
}
// Router builds the top-level HTTP handler with middleware and routes mounted.
@@ -47,6 +66,7 @@ func (h *Handler) Router() http.Handler {
r.Get("/healthz", h.Healthz)
r.Route("/api/"+APIVersion, func(r chi.Router) {
r.Use(h.rateLimit)
r.Route("/products", func(r chi.Router) {
r.Get("/barcode/{gtin}", h.ProductByBarcode)
r.Get("/search", h.SearchProducts)
+90
View File
@@ -0,0 +1,90 @@
package handler
import (
"context"
"errors"
"net"
"net/http"
"strconv"
"strings"
"time"
"github.com/baicai2026-baicai/goods/api/internal/apikey"
"github.com/baicai2026-baicai/goods/api/internal/store"
)
type ctxKey int
const apiKeyIDKey ctxKey = 0
// rateLimit authenticates an optional API key and enforces a per-minute budget
// on the public API. Anonymous callers are limited by client IP at a lower
// budget; a valid key raises the budget and attributes usage. An API key that
// is present but invalid or revoked is rejected with 401. Rate-limit headers
// are set on every response; over-budget callers get 429 + Retry-After.
func (h *Handler) rateLimit(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
id := "ip:" + clientIP(r)
limit := h.anonLimit
keyID := ""
if raw := presentedKey(r); raw != "" {
if !apikey.IsWellFormed(raw) {
writeError(w, r, http.StatusUnauthorized, "invalid_api_key", "API key 格式无效")
return
}
k, err := h.store.APIKeyByHash(r.Context(), apikey.Hash(raw))
if errors.Is(err, store.ErrNotFound) {
writeError(w, r, http.StatusUnauthorized, "invalid_api_key", "API key 无效或已吊销")
return
}
if err != nil {
writeError(w, r, http.StatusInternalServerError, "internal_error", "internal server error")
return
}
keyID = k.ID
limit = k.RateLimitPerMin
id = "key:" + k.ID
}
res := h.limiter.Allow(r.Context(), id, limit, time.Minute)
w.Header().Set("X-RateLimit-Limit", strconv.Itoa(res.Limit))
w.Header().Set("X-RateLimit-Remaining", strconv.Itoa(res.Remaining))
w.Header().Set("X-RateLimit-Reset", strconv.FormatInt(res.ResetUnix, 10))
if !res.Allowed {
retry := res.ResetUnix - time.Now().Unix()
if retry < 1 {
retry = 1
}
w.Header().Set("Retry-After", strconv.FormatInt(retry, 10))
writeError(w, r, http.StatusTooManyRequests, "rate_limited", "请求过于频繁,请稍后再试")
return
}
if keyID != "" {
h.limiter.RecordUsage(r.Context(), keyID)
next.ServeHTTP(w, r.WithContext(context.WithValue(r.Context(), apiKeyIDKey, keyID)))
return
}
next.ServeHTTP(w, r)
})
}
// presentedKey extracts an API key from the X-API-Key header or a Bearer token.
func presentedKey(r *http.Request) string {
if v := strings.TrimSpace(r.Header.Get("X-API-Key")); v != "" {
return v
}
if v := r.Header.Get("Authorization"); strings.HasPrefix(v, "Bearer ") {
return strings.TrimSpace(strings.TrimPrefix(v, "Bearer "))
}
return ""
}
// clientIP returns the caller IP, preferring chi's RealIP-normalized RemoteAddr.
func clientIP(r *http.Request) string {
if host, _, err := net.SplitHostPort(r.RemoteAddr); err == nil {
return host
}
return r.RemoteAddr
}
+120
View File
@@ -0,0 +1,120 @@
package handler
import (
"context"
"fmt"
"net/http"
"net/http/httptest"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/baicai2026-baicai/goods/api/internal/apikey"
"github.com/baicai2026-baicai/goods/api/internal/ratelimit"
"github.com/baicai2026-baicai/goods/api/internal/store"
)
// newRateLimitedHandler builds a handler backed by the test DB and a live Redis
// limiter, plus a freshly issued API key with the given per-minute limit. It
// skips when either backend is unavailable.
func newRateLimitedHandler(t *testing.T, keyLimit int) (h *Handler, plaintextKey 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 hasTable bool
if err := pool.QueryRow(ctx, "SELECT to_regclass('public.api_key') IS NOT NULL").Scan(&hasTable); err != nil || !hasTable {
pool.Close()
t.Skip("migrations not applied (api_key missing)")
}
redisURL := os.Getenv("OPENGOODS_REDIS_URL")
if redisURL == "" {
redisURL = "redis://localhost:6379/0"
}
limiter := ratelimit.New(redisURL)
pingCtx, pingCancel := context.WithTimeout(context.Background(), time.Second)
defer pingCancel()
if err := limiter.Ping(pingCtx); err != nil {
pool.Close()
t.Skipf("redis not reachable: %v", err)
}
key, hash, prefix, err := apikey.Generate()
if err != nil {
pool.Close()
t.Fatal(err)
}
name := fmt.Sprintf("test-key-%d", time.Now().UnixNano())
if _, err := pool.Exec(context.Background(),
`INSERT INTO api_key (name, key_prefix, key_hash, rate_limit_per_min) VALUES ($1,$2,$3,$4)`,
name, prefix, hash, keyLimit); err != nil {
pool.Close()
t.Fatalf("insert api_key: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), "DELETE FROM api_key WHERE key_hash=$1", hash)
pool.Close()
})
return New(store.New(pool), nil).WithRateLimit(limiter, 60), key
}
func TestRateLimitHeadersAndKeyAuth(t *testing.T) {
h, key := newRateLimitedHandler(t, 100)
req := httptest.NewRequest(http.MethodGet, "/api/"+APIVersion+"/categories", nil)
req.Header.Set("X-API-Key", key)
rec := httptest.NewRecorder()
h.Router().ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String())
}
if got := rec.Header().Get("X-RateLimit-Limit"); got != "100" {
t.Fatalf("X-RateLimit-Limit = %q, want 100 (key limit)", got)
}
if rec.Header().Get("X-RateLimit-Remaining") == "" {
t.Fatal("missing X-RateLimit-Remaining header")
}
}
func TestInvalidKeyRejected(t *testing.T) {
h, _ := newRateLimitedHandler(t, 100)
req := httptest.NewRequest(http.MethodGet, "/api/"+APIVersion+"/categories", nil)
req.Header.Set("X-API-Key", "og_live_thiskeydoesnotexist123456")
rec := httptest.NewRecorder()
h.Router().ServeHTTP(rec, req)
if rec.Code != http.StatusUnauthorized {
t.Fatalf("status = %d, want 401; body = %s", rec.Code, rec.Body.String())
}
}
func TestRateLimitExceeded(t *testing.T) {
h, key := newRateLimitedHandler(t, 1)
do := func() int {
req := httptest.NewRequest(http.MethodGet, "/api/"+APIVersion+"/categories", nil)
req.Header.Set("X-API-Key", key)
rec := httptest.NewRecorder()
h.Router().ServeHTTP(rec, req)
return rec.Code
}
if code := do(); code != http.StatusOK {
t.Fatalf("first request status = %d, want 200", code)
}
if code := do(); code != http.StatusTooManyRequests {
t.Fatalf("second request status = %d, want 429", code)
}
}
+132
View File
@@ -0,0 +1,132 @@
// Package ratelimit provides a Redis-backed fixed-window rate limiter and
// lightweight per-key usage counters for the public API.
//
// All state lives in Redis so it is shared across API replicas and visible to
// the admin console, and so the public server keeps its read-only contract
// against PostgreSQL. Every operation fails open: if Redis is unavailable the
// limiter allows the request rather than taking the API down.
package ratelimit
import (
"context"
"fmt"
"log"
"time"
"github.com/redis/go-redis/v9"
)
// Limiter throttles callers and records usage. A nil-backed Limiter (when Redis
// could not be configured) disables limiting and usage tracking.
type Limiter struct {
rdb *redis.Client
}
// Result describes the outcome of an Allow check and the headers to surface.
type Result struct {
Allowed bool
Limit int
Remaining int
ResetUnix int64
}
// UsageStat is the aggregated usage for a single API key.
type UsageStat struct {
Total int64 `json:"total"`
Today int64 `json:"today"`
LastUsedAt *int64 `json:"last_used_at,omitempty"`
}
// New builds a Limiter from a redis:// URL. On a parse error it logs and returns
// a fail-open limiter (Redis disabled) so the server still boots.
func New(redisURL string) *Limiter {
opt, err := redis.ParseURL(redisURL)
if err != nil {
log.Printf("ratelimit: invalid redis url %q: %v (rate limiting disabled)", redisURL, err)
return &Limiter{}
}
return &Limiter{rdb: redis.NewClient(opt)}
}
// Enabled reports whether a Redis backend is configured.
func (l *Limiter) Enabled() bool { return l != nil && l.rdb != nil }
// Ping verifies the Redis backend is reachable. Returns an error if disabled or
// unreachable.
func (l *Limiter) Ping(ctx context.Context) error {
if !l.Enabled() {
return redis.ErrClosed
}
return l.rdb.Ping(ctx).Err()
}
// Allow records a hit for id within a fixed window and reports whether the
// caller is under limit. Fails open (Allowed=true) on any Redis error.
func (l *Limiter) Allow(ctx context.Context, id string, limit int, window time.Duration) Result {
reset := func() int64 {
win := int64(window / time.Second)
if win < 1 {
win = 1
}
return (time.Now().Unix()/win + 1) * win
}
if !l.Enabled() {
return Result{Allowed: true, Limit: limit, Remaining: limit, ResetUnix: reset()}
}
win := int64(window / time.Second)
if win < 1 {
win = 1
}
bucket := time.Now().Unix() / win
key := fmt.Sprintf("rl:%s:%d", id, bucket)
n, err := l.rdb.Incr(ctx, key).Result()
if err != nil {
return Result{Allowed: true, Limit: limit, Remaining: limit, ResetUnix: (bucket + 1) * win}
}
if n == 1 {
l.rdb.Expire(ctx, key, time.Duration(win)*time.Second)
}
remaining := limit - int(n)
if remaining < 0 {
remaining = 0
}
return Result{
Allowed: int(n) <= limit,
Limit: limit,
Remaining: remaining,
ResetUnix: (bucket + 1) * win,
}
}
// RecordUsage increments total/daily counters and stamps last-used for a key.
// Best-effort: errors are ignored.
func (l *Limiter) RecordUsage(ctx context.Context, keyID string) {
if !l.Enabled() || keyID == "" {
return
}
now := time.Now()
day := now.Format("20060102")
pipe := l.rdb.Pipeline()
pipe.Incr(ctx, "usage:total:"+keyID)
dayKey := "usage:day:" + keyID + ":" + day
pipe.Incr(ctx, dayKey)
pipe.Expire(ctx, dayKey, 90*24*time.Hour)
pipe.Set(ctx, "usage:last:"+keyID, now.Unix(), 0)
_, _ = pipe.Exec(ctx)
}
// Usage reads aggregated usage for a key. Returns a zero-value stat on error.
func (l *Limiter) Usage(ctx context.Context, keyID string) UsageStat {
var st UsageStat
if !l.Enabled() || keyID == "" {
return st
}
day := time.Now().Format("20060102")
st.Total, _ = l.rdb.Get(ctx, "usage:total:"+keyID).Int64()
st.Today, _ = l.rdb.Get(ctx, "usage:day:"+keyID+":"+day).Int64()
if v, err := l.rdb.Get(ctx, "usage:last:"+keyID).Int64(); err == nil {
st.LastUsedAt = &v
}
return st
}
+94
View File
@@ -0,0 +1,94 @@
package ratelimit
import (
"context"
"fmt"
"os"
"testing"
"time"
)
// TestDisabledFailsOpen verifies that a Limiter without a Redis backend allows
// all requests and reports usage as zero rather than erroring.
func TestDisabledFailsOpen(t *testing.T) {
l := New("not-a-valid-url") // parse error => disabled
if l.Enabled() {
t.Fatal("expected limiter to be disabled for invalid url")
}
res := l.Allow(context.Background(), "x", 1, time.Minute)
if !res.Allowed || res.Remaining != 1 {
t.Fatalf("disabled limiter must fail open: %+v", res)
}
// Must not panic and must return zero usage.
l.RecordUsage(context.Background(), "k1")
if u := l.Usage(context.Background(), "k1"); u.Total != 0 {
t.Fatalf("disabled usage should be zero, got %+v", u)
}
}
// TestNilReceiverSafe ensures a nil *Limiter is safe to use (handler default).
func TestNilReceiverSafe(t *testing.T) {
var l *Limiter
if l.Enabled() {
t.Fatal("nil limiter must report disabled")
}
res := l.Allow(context.Background(), "x", 5, time.Minute)
if !res.Allowed {
t.Fatal("nil limiter must fail open")
}
l.RecordUsage(context.Background(), "k")
_ = l.Usage(context.Background(), "k")
}
func testLimiter(t *testing.T) *Limiter {
t.Helper()
url := os.Getenv("OPENGOODS_REDIS_URL")
if url == "" {
url = "redis://localhost:6379/0"
}
l := New(url)
if !l.Enabled() {
t.Skip("redis not configured")
}
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
if err := l.rdb.Ping(ctx).Err(); err != nil {
t.Skipf("redis not reachable: %v", err)
}
return l
}
func TestAllowFixedWindow(t *testing.T) {
l := testLimiter(t)
ctx := context.Background()
id := fmt.Sprintf("test:%d", time.Now().UnixNano())
for i := 1; i <= 2; i++ {
if res := l.Allow(ctx, id, 2, time.Minute); !res.Allowed {
t.Fatalf("request %d should be allowed: %+v", i, res)
}
}
res := l.Allow(ctx, id, 2, time.Minute)
if res.Allowed {
t.Fatalf("3rd request over limit 2 should be denied: %+v", res)
}
if res.Remaining != 0 {
t.Fatalf("remaining should be 0 when over limit, got %d", res.Remaining)
}
}
func TestRecordAndReadUsage(t *testing.T) {
l := testLimiter(t)
ctx := context.Background()
key := fmt.Sprintf("usagekey:%d", time.Now().UnixNano())
l.RecordUsage(ctx, key)
l.RecordUsage(ctx, key)
u := l.Usage(ctx, key)
if u.Total != 2 || u.Today != 2 {
t.Fatalf("expected total=2 today=2, got %+v", u)
}
if u.LastUsedAt == nil {
t.Fatal("expected last-used timestamp to be set")
}
}
+24
View File
@@ -305,6 +305,30 @@ func (s *Store) ListCategories(ctx context.Context) ([]Category, error) {
return out, rows.Err()
}
// APIKey is the minimal metadata the public API needs to authorize a caller.
type APIKey struct {
ID string
Name string
RateLimitPerMin int
}
// APIKeyByHash returns the active (non-revoked) key matching a SHA-256 hash,
// or ErrNotFound if no such active key exists.
func (s *Store) APIKeyByHash(ctx context.Context, hash string) (*APIKey, error) {
var k APIKey
err := s.pool.QueryRow(ctx,
`SELECT id, name, rate_limit_per_min
FROM api_key WHERE key_hash = $1 AND revoked_at IS NULL`, hash,
).Scan(&k.ID, &k.Name, &k.RateLimitPerMin)
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrNotFound
}
if err != nil {
return nil, err
}
return &k, nil
}
// Source describes a data source with its license and trust weight.
type Source struct {
ID string `json:"id"`
+103 -6
View File
@@ -27,6 +27,9 @@ _DEFAULT_MIN_INTERVAL = 4.0
_API_URL = "https://world.openfoodfacts.org/api/v2/product/{barcode}.json"
_SEARCH_URL = "https://world.openfoodfacts.org/api/v2/search"
# HTTP statuses worth retrying: rate limiting and transient server errors.
_RETRY_STATUS = frozenset({429, 500, 502, 503, 504})
# Fields requested from the search API so a returned product can be transformed
# without an extra per-barcode round trip.
_SEARCH_FIELDS = (
@@ -46,9 +49,13 @@ class OpenFoodFactsAdapter:
self,
client: httpx.Client | None = None,
min_interval: float = _DEFAULT_MIN_INTERVAL,
max_retries: int = 4,
backoff_base: float = 2.0,
) -> None:
self._client = client or httpx.Client(headers={"User-Agent": USER_AGENT}, timeout=30.0)
self._min_interval = min_interval
self._max_retries = max_retries
self._backoff_base = backoff_base
self._last_call = 0.0
def _throttle(self) -> None:
@@ -58,11 +65,49 @@ class OpenFoodFactsAdapter:
time.sleep(wait)
self._last_call = time.monotonic()
def _get(self, url: str, params: dict | None = None) -> httpx.Response:
"""GET with throttling and retry/backoff on transient errors.
Retries on connection/timeout errors and on retryable HTTP statuses
(429 and 5xx, which OFF returns intermittently when overloaded), using
exponential backoff that honours a ``Retry-After`` header when present.
"""
last_exc: Exception | None = None
for attempt in range(self._max_retries + 1):
self._throttle()
try:
resp = self._client.get(url, params=params)
except httpx.TransportError as exc:
last_exc = exc
else:
if resp.status_code < 400 or resp.status_code not in _RETRY_STATUS:
resp.raise_for_status()
return resp
last_exc = httpx.HTTPStatusError(
f"retryable status {resp.status_code}", request=resp.request, response=resp
)
if attempt < self._max_retries:
retry_after = self._retry_after(last_exc)
time.sleep(retry_after if retry_after is not None else self._backoff_base**attempt)
assert last_exc is not None
raise last_exc
@staticmethod
def _retry_after(exc: Exception | None) -> float | None:
resp = getattr(exc, "response", None)
if resp is None:
return None
value = resp.headers.get("Retry-After")
if not value:
return None
try:
return float(value)
except ValueError:
return None
def fetch_barcode(self, barcode: str) -> dict | None:
"""Fetch a single product by barcode; return the raw `product` dict."""
self._throttle()
resp = self._client.get(_API_URL.format(barcode=barcode))
resp.raise_for_status()
resp = self._get(_API_URL.format(barcode=barcode))
payload = resp.json()
if payload.get("status") != 1:
return None
@@ -91,8 +136,7 @@ class OpenFoodFactsAdapter:
``last_modified_t`` they processed as the next watermark.
"""
for page in range(1, max_pages + 1):
self._throttle()
resp = self._client.get(
resp = self._get(
_SEARCH_URL,
params={
"fields": _SEARCH_FIELDS,
@@ -101,7 +145,6 @@ class OpenFoodFactsAdapter:
"page_size": page_size,
},
)
resp.raise_for_status()
products = resp.json().get("products") or []
if not products:
return
@@ -114,6 +157,60 @@ class OpenFoodFactsAdapter:
if reached_old or len(products) < page_size:
return
def fetch_by_country(
self,
country: str,
*,
page_size: int = 100,
max_pages: int = 10,
sort_by: str = "unique_scans_n",
) -> Iterator[dict]:
"""Yield products sold in ``country`` (an OFF ``countries_tags_en`` slug).
Used to seed a market-specific catalogue (e.g. ``china``). Results are
sorted by ``sort_by`` (default ``unique_scans_n`` so the most-scanned,
best-known products come first) and de-duplicated across pages, since
OFF's popularity ordering is not stable between page requests.
"""
seen: set[str] = set()
for page in range(1, max_pages + 1):
resp = self._get(
_SEARCH_URL,
params={
"fields": _SEARCH_FIELDS,
"countries_tags_en": country,
"sort_by": sort_by,
"page": page,
"page_size": page_size,
},
)
products = resp.json().get("products") or []
if not products:
return
new_on_page = 0
for prod in products:
code = str(prod.get("code") or "")
if code and code in seen:
continue
if code:
seen.add(code)
new_on_page += 1
yield prod
if len(products) < page_size or new_on_page == 0:
return
def is_cn_gs1(code: str | None) -> bool:
"""Return True for a GS1 China company prefix (barcodes starting 690-699).
These identify products registered with GS1 China, i.e. genuinely domestic
items, as opposed to imported goods merely tagged as sold in China.
"""
if not code:
return False
code = code.strip()
return len(code) >= 3 and code[:2] == "69" and code[2].isdigit()
def read_dump(path: str | Path) -> Iterator[dict]:
"""Yield raw product records from an OFF JSONL dump file.
+22
View File
@@ -7,6 +7,7 @@ source with field-level provenance in `product_source`.
from __future__ import annotations
import json
import logging
import os
from typing import Any
@@ -18,6 +19,8 @@ from opengoods.etl.quality import update_quality
OFF_HOMEPAGE = "https://world.openfoodfacts.org"
logger = logging.getLogger(__name__)
def default_dsn() -> str:
return os.environ.get(
@@ -196,6 +199,25 @@ def load_record(conn: psycopg.Connection, rec: dict[str, Any], source_id: str, r
return product_id
def load_record_safe(
conn: psycopg.Connection, rec: dict[str, Any], source_id: str, raw: dict
) -> bool:
"""Load one record inside a savepoint.
On success the record's writes stay in the surrounding transaction. On any
error, only this record's writes are rolled back (to the savepoint) and the
batch continues, so a single malformed source record cannot abort a large
import. Returns True if loaded, False if skipped due to an error.
"""
try:
with conn.transaction():
load_record(conn, rec, source_id, raw)
return True
except Exception as exc: # noqa: BLE001 - per-record isolation is intentional
logger.warning("skipping record gtin=%s: %s", rec.get("gtin"), exc)
return False
def _jsonable(raw: dict) -> dict:
"""Drop values that are not JSON-serializable from a raw record."""
try:
+11 -3
View File
@@ -78,6 +78,14 @@ def map_category(raw: dict) -> str | None:
return None
def _clamp(value: str | None, max_len: int) -> str | None:
"""Trim a string to fit a bounded DB column; external data length varies."""
if value is None:
return None
value = value.strip()
return value[:max_len] or None
def _clean_tags(tags: list[str] | None, prefix: str = "") -> list[str]:
out: list[str] = []
for t in tags or []:
@@ -143,16 +151,16 @@ def transform(raw: dict) -> dict | None:
"brand": brand,
"category_path": map_category(raw),
"net_content_value": net_value,
"net_content_unit": net_unit,
"net_content_unit": _clamp(net_unit, 16),
"net_content_canonical": net_canonical,
"country_of_origin": (raw.get("countries") or "").split(",")[0].strip() or None,
"country_of_origin": _clamp((raw.get("countries") or "").split(",")[0].strip() or None, 64),
"food": {
"ingredients_text": raw.get("ingredients_text") or None,
"allergens": _clean_tags(raw.get("allergens_tags")),
"additives": _clean_tags(raw.get("additives_tags")),
"nutriments": transform_nutriments(raw.get("nutriments") or {}),
"nutrition_basis": "per_100g",
"serving_size": raw.get("serving_size") or None,
"serving_size": _clamp(raw.get("serving_size") or None, 32),
"nutri_score": (raw.get("nutriscore_grade") or "").upper()[:1] or None,
},
"image_url": raw.get("image_front_url") or raw.get("image_url") or None,
+40 -9
View File
@@ -7,6 +7,11 @@ Usage:
# from a downloaded OFF JSONL dump (optionally .gz), limited to N records
python -m opengoods.jobs.seed_off --dump products.jsonl.gz --limit 1000
# market-focused: the most-scanned products sold in China, restricted to
# genuine GS1-China (69x) barcodes
python -m opengoods.jobs.seed_off --country china --domestic-only \
--max-pages 20 --limit 1000
The OFF read API is rate-limited client-side; for large imports use a dump.
"""
@@ -18,25 +23,38 @@ from collections.abc import Iterator
import psycopg
from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter, read_dump
from opengoods.etl.load import default_dsn, ensure_source, load_record
from opengoods.adapters.openfoodfacts import (
OpenFoodFactsAdapter,
is_cn_gs1,
read_dump,
)
from opengoods.etl.load import default_dsn, ensure_source, load_record_safe
from opengoods.etl.transform import transform
def _raw_records(args: argparse.Namespace) -> Iterator[dict]:
if args.dump:
records = read_dump(args.dump)
records: Iterator[dict] = read_dump(args.dump)
elif args.country:
adapter = OpenFoodFactsAdapter(min_interval=args.min_interval)
records = adapter.fetch_by_country(
args.country, page_size=args.page_size, max_pages=args.max_pages
)
else:
adapter = OpenFoodFactsAdapter(min_interval=args.min_interval)
records = adapter.fetch(args.barcodes)
for i, rec in enumerate(records):
if args.limit and i >= args.limit:
yielded = 0
for rec in records:
if args.domestic_only and not is_cn_gs1(str(rec.get("code") or "")):
continue
if args.limit and yielded >= args.limit:
break
yielded += 1
yield rec
def run(args: argparse.Namespace) -> int:
loaded = skipped = 0
loaded = skipped = errored = 0
with psycopg.connect(args.dsn, autocommit=False) as conn:
source_id = ensure_source(conn)
for raw in _raw_records(args):
@@ -44,10 +62,12 @@ def run(args: argparse.Namespace) -> int:
if rec is None:
skipped += 1
continue
load_record(conn, rec, source_id, raw)
loaded += 1
if load_record_safe(conn, rec, source_id, raw):
loaded += 1
else:
errored += 1
conn.commit()
print(f"loaded={loaded} skipped={skipped}")
print(f"loaded={loaded} skipped={skipped} errored={errored}")
return 0
@@ -56,7 +76,18 @@ def main(argv: list[str] | None = None) -> int:
src = parser.add_mutually_exclusive_group(required=True)
src.add_argument("--barcodes", nargs="+", help="barcodes to fetch via the OFF API")
src.add_argument("--dump", help="path to an OFF JSONL dump (.jsonl or .jsonl.gz)")
src.add_argument(
"--country",
help="OFF countries_tags_en slug to seed from, e.g. 'china' (most-scanned first)",
)
parser.add_argument(
"--domestic-only",
action="store_true",
help="keep only genuine GS1-China (69x) barcodes; drop imported goods",
)
parser.add_argument("--limit", type=int, default=0, help="max records to load (0 = all)")
parser.add_argument("--page-size", type=int, default=100, help="search page size")
parser.add_argument("--max-pages", type=int, default=10, help="max search pages (country mode)")
parser.add_argument("--min-interval", type=float, default=4.0, help="API throttle seconds")
parser.add_argument("--dsn", default=default_dsn(), help="PostgreSQL DSN")
return run(parser.parse_args(argv))
+11 -6
View File
@@ -17,14 +17,14 @@ import sys
import psycopg
from opengoods.adapters.openfoodfacts import SOURCE_NAME, OpenFoodFactsAdapter
from opengoods.etl.load import default_dsn, ensure_source, load_record
from opengoods.etl.load import default_dsn, ensure_source, load_record_safe
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
loaded = skipped = errored = 0
high_watermark = 0
with psycopg.connect(args.dsn, autocommit=False) as conn:
source_id = ensure_source(conn)
@@ -38,16 +38,21 @@ def run(args: argparse.Namespace) -> int:
if rec is None:
skipped += 1
continue
load_record(conn, rec, source_id, raw)
loaded += 1
if load_record_safe(conn, rec, source_id, raw):
loaded += 1
else:
errored += 1
set_watermark(
conn,
SOURCE_NAME,
high_watermark,
stats={"loaded": loaded, "skipped": skipped, "since": since},
stats={"loaded": loaded, "skipped": skipped, "errored": errored, "since": since},
)
conn.commit()
print(f"since={since} loaded={loaded} skipped={skipped} watermark={high_watermark}")
print(
f"since={since} loaded={loaded} skipped={skipped} "
f"errored={errored} watermark={high_watermark}"
)
return 0
+23 -1
View File
@@ -11,7 +11,7 @@ from pathlib import Path
import pytest
from opengoods.etl.load import default_dsn, ensure_source, load_record
from opengoods.etl.load import default_dsn, ensure_source, load_record, load_record_safe
from opengoods.etl.transform import transform
psycopg = pytest.importorskip("psycopg")
@@ -59,3 +59,25 @@ def test_load_record_roundtrip(conn):
assert prov[0] >= 1
conn.rollback() # keep the test DB clean
def test_load_record_safe_isolates_bad_record(conn):
source_id = ensure_source(conn)
# Unique gtin so the good record is a fresh INSERT, not an upsert/update.
good = transform(FIXTURE)
good["gtin"] = "4006381333931"
assert load_record_safe(conn, good, source_id, FIXTURE) is True
after_good = conn.execute("SELECT count(*) FROM product").fetchone()[0]
# A record whose name violates NOT NULL must not abort the batch.
bad = dict(good)
bad["gtin"] = "5000112637922"
bad["name"] = None
assert load_record_safe(conn, bad, source_id, {}) is False
# The good record survived the bad one's rollback-to-savepoint.
after_bad = conn.execute("SELECT count(*) FROM product").fetchone()[0]
assert after_bad == after_good
conn.rollback() # keep the test DB clean
+53
View File
@@ -0,0 +1,53 @@
import httpx
from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter, is_cn_gs1
def _adapter(handler, **kwargs):
client = httpx.Client(transport=httpx.MockTransport(handler))
return OpenFoodFactsAdapter(client=client, min_interval=0, backoff_base=0, **kwargs)
def test_is_cn_gs1():
assert is_cn_gs1("6901234567892") # GS1 China prefix
assert is_cn_gs1("690")
assert not is_cn_gs1("3017624010701") # France
assert not is_cn_gs1("5449000000996") # Belgium
assert not is_cn_gs1("")
assert not is_cn_gs1(None)
assert not is_cn_gs1("69") # too short to carry a prefix digit
def test_fetch_by_country_passes_filter_and_paginates():
seen_params = []
def handler(request: httpx.Request) -> httpx.Response:
seen_params.append(dict(request.url.params))
page = int(request.url.params.get("page", "1"))
if page == 1:
return httpx.Response(
200,
json={"products": [{"code": "6901"}, {"code": "6902"}]},
)
return httpx.Response(200, json={"products": []})
adapter = _adapter(handler)
out = list(adapter.fetch_by_country("china", page_size=2, max_pages=5))
assert [p["code"] for p in out] == ["6901", "6902"]
assert seen_params[0]["countries_tags_en"] == "china"
assert seen_params[0]["sort_by"] == "unique_scans_n"
def test_fetch_by_country_dedupes_across_pages():
def handler(request: httpx.Request) -> httpx.Response:
page = int(request.url.params.get("page", "1"))
if page == 1:
return httpx.Response(200, json={"products": [{"code": "6901"}, {"code": "6902"}]})
if page == 2:
# OFF popularity ordering is unstable; a repeat appears on page 2
return httpx.Response(200, json={"products": [{"code": "6902"}, {"code": "6903"}]})
return httpx.Response(200, json={"products": []})
adapter = _adapter(handler)
out = [p["code"] for p in adapter.fetch_by_country("china", page_size=2, max_pages=5)]
assert out == ["6901", "6902", "6903"]
+59
View File
@@ -0,0 +1,59 @@
import httpx
import pytest
from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter
def _adapter(handler, **kwargs):
client = httpx.Client(transport=httpx.MockTransport(handler))
return OpenFoodFactsAdapter(client=client, min_interval=0, backoff_base=0, **kwargs)
def test_retries_transient_5xx_then_succeeds():
calls = {"n": 0}
def handler(request: httpx.Request) -> httpx.Response:
calls["n"] += 1
if calls["n"] < 3:
return httpx.Response(503)
return httpx.Response(200, json={"status": 1, "product": {"code": "x"}})
adapter = _adapter(handler, max_retries=4)
assert adapter.fetch_barcode("x") == {"code": "x"}
assert calls["n"] == 3 # two 503s retried, third succeeds
def test_retries_exhausted_raises():
def handler(request: httpx.Request) -> httpx.Response:
return httpx.Response(503)
adapter = _adapter(handler, max_retries=2)
with pytest.raises(httpx.HTTPStatusError):
adapter.fetch_barcode("x")
def test_non_retryable_4xx_not_retried():
calls = {"n": 0}
def handler(request: httpx.Request) -> httpx.Response:
calls["n"] += 1
return httpx.Response(404)
adapter = _adapter(handler, max_retries=4)
with pytest.raises(httpx.HTTPStatusError):
adapter.fetch_barcode("x")
assert calls["n"] == 1 # 404 is not retried
def test_retries_connection_error_then_succeeds():
calls = {"n": 0}
def handler(request: httpx.Request) -> httpx.Response:
calls["n"] += 1
if calls["n"] == 1:
raise httpx.ConnectError("boom")
return httpx.Response(200, json={"status": 1, "product": {"code": "y"}})
adapter = _adapter(handler, max_retries=4)
assert adapter.fetch_barcode("y") == {"code": "y"}
assert calls["n"] == 2
+11
View File
@@ -65,3 +65,14 @@ def test_transform_full_record():
def test_transform_drops_unnamed():
assert transform({"code": "0000000000000"}) is None
def test_transform_clamps_long_serving_size():
raw = {
"code": "3017624010701",
"product_name": "X",
"serving_size": "1 portion (30 g) / 1 portion (30 g) / extra long descriptive text here",
}
rec = transform(raw)
assert rec is not None
assert len(rec["food"]["serving_size"]) == 32
+1
View File
@@ -0,0 +1 @@
DROP TABLE IF EXISTS api_key;
+24
View File
@@ -0,0 +1,24 @@
-- API keys for the public read-only API. Keys grant higher rate limits and let
-- usage be attributed to a caller; the API itself stays free and read-only.
-- Only the SHA-256 hash of a key is stored; the plaintext is shown once at
-- creation time. Keys are issued/revoked from the admin console. The public
-- server only ever SELECTs from this table (request counting lives in Redis),
-- preserving its read-only contract against PostgreSQL.
CREATE TABLE IF NOT EXISTS api_key (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
name TEXT NOT NULL,
key_prefix VARCHAR(20) NOT NULL, -- shown for identification, e.g. og_live_AbC1
key_hash TEXT NOT NULL UNIQUE, -- hex SHA-256 of the full key
owner_email TEXT,
tier VARCHAR(16) NOT NULL DEFAULT 'free',
rate_limit_per_min INT NOT NULL DEFAULT 120,
revoked_at TIMESTAMPTZ,
created_by TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
CONSTRAINT api_key_tier_chk CHECK (tier IN ('free', 'partner', 'internal')),
CONSTRAINT api_key_rate_chk CHECK (rate_limit_per_min > 0)
);
-- Fast lookup of active keys by their hash on every authenticated request.
CREATE INDEX IF NOT EXISTS idx_api_key_active_hash
ON api_key (key_hash) WHERE revoked_at IS NULL;