Appearance
🏗️ GROUPBUY / PRE-ORDER PLATFORM — SYSTEM ARCHITECTURE
Phiên bản: v3.0 · Mô hình: Cọc trước → Chốt giá khi đóng nhóm → Thu phần còn lại → Vận chuyển Kèm theo: Pipeline auto-crawl tạo campaign · Quan sát hệ thống bằng OpenTelemetry Nguyên tắc stack: Miễn phí / OSS, chạy được trọn bộ trên 1 laptop bằng docker-compose
MỤC LỤC
- 0 · Tóm tắt điều hành
- 1 · Quyết định thiết kế then chốt
- 2 · Kiến trúc tổng thể
- 3 · Thẻ dịch vụ & tech stack
- 4 · Layer OpenTelemetry
- 5 · Luồng đặt slot & thanh toán cọc
- 6 · Auto-crawl pipeline
- 7 · Tiền: chốt giá · settlement · hoàn tiền
- 8 · Data model đầy đủ
- 9 · Business knobs
- 10 · Roadmap triển khai
- 11 · Chạy local: docker-compose & repo layout
- 12 · Lộ trình học
0 · TÓM TẮT ĐIỀU HÀNH
| Hạng mục | Quyết định |
|---|---|
| Sản phẩm | Nền tảng groupbuy/pre-order: creator mở campaign, user đặt cọc giữ slot, hệ thống chốt giá theo MOQ tier khi đóng nhóm, user trả phần còn lại để xác nhận và nhận hàng |
| Bất biến kinh doanh | final_price ≥ deposit luôn đúng ⇒ hệ thống chỉ thu thêm, không bao giờ hoàn chênh lệch |
| Hot path | Go Order Engine — chống oversell bằng atomic SQL trên PostgreSQL |
| Consistency | PostgreSQL là source of truth duy nhất; mọi state transition đi qua Outbox Pattern |
| Tiền | Webhook idempotent · settlement batch resumable · hoàn cọc tự động khi user bỏ qua hạn |
| Tự động hoá nguồn hàng | Crawler (Playwright) → LLM extract → Dedup → Draft PENDING_REVIEW → Approve → Publish |
| Quan sát | OpenTelemetry xuyên suốt (kể cả job delay 15 phút vẫn cùng traceId) |
| Chi phí hạ tầng học/dev | $0 |
10 quyết định kiến trúc
- PG là source of truth duy nhất cho inventory. Redis chỉ hiển thị counter (SSE) và làm queue — không bao giờ ra quyết định ghi.
- Outbox Pattern: mỗi transition ghi state + event trong cùng 1 transaction; side effects chỉ dispatch qua event/job.
- Webhook idempotency-first: verify HMAC → insert với unique key → trả 200 ngay → xử lý async.
- Chỉ đơn
DEPOSIT_HELDmới đếm vào MOQ/tier ⇒ tổngqty ổn định tuyệt đối tại thời điểm close. deposit_per_unit ≤ MIN(tier.unit_price)⇒ mọi kịch bản chốt đềufinal ≥ cọc.- Settlement 1 chiều (chỉ thu thêm) — xoá khỏi scope nhánh hoàn-chênh-lệch.
- Quá hạn settlement → tự huỷ + hoàn cọc qua RefundEngine (pattern resumable dùng lại).
- Crawler không bao giờ auto-publish 100%: luôn qua
PENDING_REVIEW; auto-publish chỉ theo confidence threshold + global kill-switch. - Observability từ ngày 1, chạy ngang mọi layer (cross-cutting).
- Media crawl luôn tải về CDN của mình — không hotlink.
2 · KIẾN TRÚC TỔNG THỂ
👤 SHOPPER 👤 CREATOR / ADMIN
│ │
│ ① HTTPS │ ① HTTPS
▼ ▼
┌───────────────────────────────────────────────────────────────┐
│ ② FRONTEND · Next.js 14 (App Router) + Cloudflare CDN │
│ Storefront(ISR) · Checkout(SSE) · AdminDash · DraftReview │
└───────┬───────────────────────────────┬───────────────────────┘
│ ② REST + SSE │ ③ GraphQL
▼ ▼
┌───────────────────────────────────────────────────────────────┐
│ ③ GATEWAY · NestJS + Fastify │
│ JWT(jose)·RBAC · Zod · Mercurius GraphQL · SSE Hub │
│ Webhook Receiver (verify HMAC → idempotency → enqueue) │
└───┬──────────────┬───────────────────┬────────────────┬────────┘
│④ gRPC │⑤ enqueue │⑥ outbox events │⑤ enqueue
▼ ▼ ▼ ▼
┌──────────┐ ┌──────────────────┐ ┌───────────────┐ ┌──────────────┐
│④ GO ORDER│ │⑤ WORKER FLEET │ │⑥ ORCHESTRATOR │ │⑦ CRAWLER + AI│
│ ENGINE │ │ DepositExpiry │ │ (State Machine│ │ PIPELINE │
│ │ │ MoqFinalizer │ │ + Outbox) │ └──────┬───────┘
│ │ │ SettlementBatch★ │ └───────┬───────┘ │
│ │ │ RefundEngine │ │ │
│ │ │ MediaPipeline │ │ │
│ │ │ Notifier │ │ │
└────┬─────┘ └──────┬───────────┘ │ │
│ │ ⑧ Stripe/VNPay · Telegram · Resend │ ⑧ LLM
▼ ▼ ▼
┌───────────────────────────────────────────────────────────────┐
│ ⑨ DATA LAYER │
│ PostgreSQL 16 + PgBouncer │ Valkey/Redis │ MinIO / R2 │
└───────────────────────────────────────────────────────────────┘
═══ TẤT CẢ SERVICE ĐỀU PHÁT TELEMETRY (OTLP) ═══
│
▼
┌─────────────────────────┐
│ ⑩ OPEN TELEMETRY COLLECTOR │
└────┬─────┬──────┬───────┘
▼ ▼ ▼
Jaeger Prom Loki ──▶ Grafana (1 màn hình duy nhất)Legend các luồng
| # | Luồng | Giao thức |
|---|---|---|
| ① | Client → Next.js | HTTPS |
| ② | Storefront → Gateway | REST JSON + SSE stream |
| ③ | Admin → Gateway | GraphQL (Mercurius) |
| ④ | Gateway ↔ Go Engine | gRPC + Protobuf |
| ⑤ | Mọi service → hàng đợi | BullMQ (Valkey/Redis) |
| ⑥ | Outbox relay → Orchestrator | Poll PG 1s (sau nâng cấp Debezium) |
| ⑦ | Cron scheduler → Crawler | BullMQ repeatable jobs |
| ⑧ | Ra thế giới ngoài | HTTPS API (Stripe/VNPay, Telegram, LLM…) |
| ⑨ | Services → DB | TCP (pgx / ioredis / S3 API) |
| ⑩ | Services → Collector | OTLP gRPC :4317 / HTTP :4318 |
3 · THẺ DỊCH VỤ & TECH STACK
┌────────────────────────────────────────────────────────────┐
│ ② FRONTEND — Next.js │
├────────────────────────────────────────────────────────────┤
│ Framework : Next.js 14 App Router + TailwindCSS │
│ State : TanStack Query (REST) + urql/Apollo (admin) │
│ Realtime : EventSource API (SSE native browser) │
│ Deploy : Cloudflare Pages (free, unlimited bandwidth) │
│ Telemetry : fetch tự gắn header W3C traceparent │
└────────────────────────────────────────────────────────────┘
┌────────────────────────────────────────────────────────────┐
│ ③ GATEWAY — NestJS │
├────────────────────────────────────────────────────────────┤
│ Runtime : Node 20 LTS + NestJS + Fastify adapter │
│ GraphQL : Mercurius (native Fastify, nhẹ hơn Apollo) │
│ Auth : jose (JWT rotation) + refresh token ở Redis │
│ Validation : Zod — share schema với FE qua monorepo pkg │
│ Rate limit : @fastify/rate-limit (Redis store) │
│ Telemetry : @opentelemetry/auto-instrumentations-node │
└────────────────────────────────────────────────────────────┘
┌────────────────────────────────────────────────────────────┐
│ ④ GO ORDER ENGINE │
├────────────────────────────────────────────────────────────┤
│ Runtime : Go 1.22 │
│ RPC : google.golang.org/grpc + protobuf │
│ DB Driver : pgx/v5 (raw pool, KHÔNG ORM cho hot path) │
│ Codegen : sqlc (SQL → typed Go) │
│ Contract : packages/proto — .proto dùng chung │
│ Telemetry : otelgrpc (interceptor) + otelpgx (DB spans) │
│ Image : distroless/alpine ~20MB │
└────────────────────────────────────────────────────────────┘
┌────────────────────────────────────────────────────────────┐
│ ⑥ ORCHESTRATOR (trong Gateway hoặc service riêng) │
├────────────────────────────────────────────────────────────┤
│ State machine : xstate (vẽ được chart trực quan) │
│ Outbox relay : poll bảng outbox mỗi 1s (đủ cho học) │
│ Telemetry : span riêng cho mỗi transition │
└────────────────────────────────────────────────────────────┘
┌────────────────────────────────────────────────────────────┐
│ ⑤ WORKER FLEET │
├────────────────────────────────────────────────────────────┤
│ Queue : BullMQ + ioredis (delayed, repeatable, DLQ, │
│ retry backoff — có sẵn hết) │
│ Workers : DepositExpiry · MoqFinalizer · SettlementBatch│
│ RefundEngine · MediaPipeline(sharp) · Notifier│
│ Payment : Stripe SDK (test mode) / VNPay sandbox │
│ Notify : Telegraf (Telegram Bot) + Resend (email free) │
│ Telemetry : ⚠️ inject/extract trace context vào job.data │
└────────────────────────────────────────────────────────────┘
┌────────────────────────────────────────────────────────────┐
│ ⑦ CRAWLER + AI PIPELINE │
├────────────────────────────────────────────────────────────┤
│ Crawl : Playwright (headless Chromium) + Cheerio │
│ Schedule : BullMQ repeatable jobs │
│ Storage : raw HTML/media → MinIO (docker) hoặc CF R2 │
│ LLM : Ollama local (qwen2.5:7b — FREE 100%) │
│ fallback cloud: Groq / Gemini API free tier │
│ AI SDK : Vercel AI SDK (đổi provider dễ) │
│ Schema : Zod parse JSON LLM → fail → review queue │
│ Dedup : URL hash + cosine similarity (transformers.js)│
│ Alert : Telegraf → Telegram │
└────────────────────────────────────────────────────────────┘
┌────────────────────────────────────────────────────────────┐
│ ⑨ DATA LAYER │
├────────────────────────────────────────────────────────────┤
│ OLTP : PostgreSQL 16 + PgBouncer (transaction mode) │
│ Host học : Docker local, hoặc Neon free tier │
│ Cache/Queue: Valkey 7 (BSD) hoặc Redis 7 │
│ Objects : MinIO (S3-compatible) │
│ Migration : Drizzle Kit hoặc node-pg-migrate │
└────────────────────────────────────────────────────────────┘Bảng chi phí — toàn bộ $0
| Thành phần | Chọn | Giấy phép / Free tier |
|---|---|---|
| Frontend host | Cloudflare Pages | Free unlimited BW |
| API framework | NestJS + Fastify + Mercurius | MIT |
| Auth | jose + Redis | MIT |
| Hot path | Go + grpc + pgx + sqlc | BSD/MIT |
| Queue | BullMQ + Valkey (hoặc Redis 7) | MIT / BSD |
| Database | Postgres 16 + PgBouncer (Docker hoặc Neon free) | OSS |
| Object storage | MinIO local / Cloudflare R2 (10GB free) | AGPL / free tier |
| Crawler | Playwright + Cheerio | Apache-2.0 / MIT |
| LLM | Ollama + qwen2.5:7b (local) / Groq free tier | MIT / free |
| Payment test | Stripe test mode / VNPay sandbox | Free |
| Notify | Telegram Bot API + Resend (3k email/tháng) | Free |
| Tracing | OTel SDK + Collector + Jaeger | Apache-2.0/CNCF |
| Metrics/Dashboard | Prometheus + Grafana | Apache-2.0/AGPL |
| Logs | pino (+ Loki nếu muốn tập trung) | MIT/AGPL |
| Error tracking | GlitchTip self-host (hoặc Sentry free 5k err) | MIT/BSD |
4 · LAYER OPEN TELEMETRY
4.1 Lan truyền trace — một request, một traceId xuyên suốt
POST /checkout (browser gắn header: traceparent: 00-{trace-id}-{span-id}-01)
│
▼
┌─ WATERFALL NHÌN ĐƯỢC TRONG JAEGER ─────────────────────────────┐
│ ▓ gateway: POST /orders 128ms │
│ ├─▓ auth.verify_jwt 1ms │
│ ├─▓ zod.validate_body 1ms │
│ ├─▓ grpc.ReserveSlot ──── metadata.traceparent ──▶ 19ms │
│ │ ▓ go-engine.ReserveSlot 18ms │
│ │ ▓ pgx.tx: UPDATE campaigns … 13ms │
│ ├─▓ bullmq.enqueue(expiry:order:991) 2ms │
│ └─▓ redis.publish → SSE fan-out 1ms │
└────────────────────────────────────────────────────────────────┘
│
│ ... 15 PHÚT SAU — VẪN CÙNG traceId ...
▼
▓ worker.processJob(expiry:order:991) ← extract context từ job.data
└─▓ pgx.tx: release expired slots4.2 Instrumentation từng service
| Service | SDK / Lib | Ghi chú |
|---|---|---|
| NestJS Gateway | @opentelemetry/sdk-node + auto-instrumentations-node | Auto bắt http, graphql, ioredis, pg |
| Worker Fleet | Cùng SDK + inject thủ công qua job.data.__otel | BullMQ không tự propagate |
| Go Engine | go.opentelemetry.io/otel + otelgrpc + otelpgx | Interceptor gRPC 2 phía |
| Frontend | Browser tự gửi header traceparent | Không cần SDK nặng |
| Collector | otelcol-contrib (docker image chính thức) | Nhận OTLP :4317/:4318 |
4.3 Inject/extract trace context qua BullMQ (bắt buộc tự viết)
ts
// PRODUCER — khi enqueue
import { context, propagation } from '@opentelemetry/api';
const carrier: Record<string, string> = {};
propagation.inject(context.active(), carrier);
await queue.add('expiry', { orderId, __otel: carrier }, { delay: 15 * 60_000 });
// CONSUMER — khi xử lý job
const ctx = propagation.extract(context.active(), job.data.__otel);
await context.with(ctx, () => processExpiry(job)); // span nối tiếp cùng trace4.4 Cấu hình Collector tối thiểu (infra/otel/collector.yaml)
yaml
receivers:
otlp:
protocols:
grpc: { endpoint: 0.0.0.0:4317 }
http: { endpoint: 0.0.0.0:4318 }
processors:
memory_limiter: { check_interval: 1s, limit_mib: 256 }
batch: { timeout: 5s }
exporters:
otlp/jaeger:
endpoint: jaeger:4317
tls: { insecure: true }
prometheus:
endpoint: 0.0.0.0:8889
service:
pipelines:
traces: { receivers: [otlp], processors: [memory_limiter, batch], exporters: [otlp/jaeger] }
metrics: { receivers: [otlp], processors: [batch], exporters: [prometheus] }4.5 Sampling
ts
// dev: 100% để học; prod: 10% đủ đọc xu hướng
sampler: new ParentBasedSampler(
new TraceIdRatioBased(process.env.NODE_ENV === 'production' ? 0.1 : 1.0),
);5 · LUỒNG ĐẶT SLOT & THANH TOÁN CỌC
5.1 Flow A — Đặt slot (hot path, chống oversell tuyệt đối)
sql
-- Go Engine, 1 transaction duy nhất:
BEGIN;
UPDATE campaigns
SET taken = taken + :qty
WHERE id = :campaign_id
AND taken + :qty <= capacity
RETURNING capacity - taken AS remaining; -- 0 rows ⟹ HẾT SLOT, rollback
INSERT INTO orders (status, expires_at, deposit_due)
VALUES ('RESERVED', now() + interval '15 minutes', :deposit);
INSERT INTO outbox (event_type, payload) VALUES ('order.reserved', ...);
COMMIT;Dạng
UPDATE … WHERE taken + qty <= capacity RETURNINGlà atomic single-statement — nhanh hơnSELECT FOR UPDATEvì không giữ lock chờ nhau. Đây là vũ khí chính chống oversell khi 500 người cùng giành 1 slot cuối.
Sau commit: outbox relay đẩy event → Redis Pub/Sub → SSE cập nhật counter mọi client → enqueue deposit-expiry job (delay 15 phút).
5.2 Flow B — Reservation 15 phút (không ghost-hold)
BullMQ delayed job (jobId = "expiry:order:{id}", delay = 15m)
│ wake-up
▼
RE-CHECK trong 1 transaction (có row lock):
order còn 'RESERVED'? payment đã confirm?
├─ đã trả tiền → exit im lặng (idempotent)
└─ chưa trả → UPDATE orders SET status='EXPIRED'
UPDATE slots SET taken = taken - :qty
INSERT outbox('order.expired')Race case: user trả tiền ở phút 14:59 → webhook (Flow C) và expiry job đụng nhau → cả hai đều phải FOR UPDATE row order trước khi quyết định → ai vào trước thắng.
5.3 Flow C — Payment webhook (idempotent bắt buộc)
Provider ─▶ POST /webhooks/payments
1. Verify HMAC signature → sai: 401
2. INSERT INTO webhook_events (idempotency_key UNIQUE, payload)
ON CONFLICT (idempotency_key) DO NOTHING;
3. rowcount = 0 → trả 200 (trùng, đã xử lý rồi) ✅
rowcount = 1 → trả 200 NGAY + enqueue ProcessPayment
4. Worker ProcessPayment (tx riêng):
lock order → mark PAID → INSERT outbox('payment.confirmed')6 · AUTO-CRAWL PIPELINE
6.1 Flow D — Crawl → Extract → Draft → Publish
[Cron Scheduler — lịch riêng từng domain]
▼
Crawler Service
• Playwright headless pool + proxy rotator
• Per-domain rate limit + robots.txt + backoff + captcha flag
▼
Raw Store (S3/MinIO): html.gz + screenshot + metadata ← luôn lưu gốc
▼
LLM Extraction Worker
raw HTML ──prompt──▶ JSON draft
→ Zod validate (giá, MOQ, tồn kho, lead time…)
→ confidence score
├─ hợp lệ + conf đủ cao ──────────▶ tiếp tục
└─ lỗi schema / conf thấp ──▶ REVIEW_QUEUE (xem raw + draft cạnh nhau)
▼
Dedup: canonical URL hash + embedding similarity tên SP
├─ trùng campaign đang OPEN ──▶ route sang PriceChangeDetector (6.2)
└─ mới ──▶ tải media về CDN mình + tạo CampaignDraft(PENDING_REVIEW)
▼
Admin Draft Review UI: chỉnh giá/MOQ/tier/deposit → APPROVE
▼
Orchestrator.transition(DRAFT → OPEN) ← đi qua cùng state machine,
KHÔNG có đường tắt riêngNguyên tắc an toàn:
- Giai đoạn đầu: human review 100%. Auto-publish chỉ bật khi
confidence ≥ 0.95và source nằm trong whitelist tin cậy, luôn kèm global kill-switch + audit log. - Media luôn download về CDN mình (link chết giữa campaign = thảm họa).
- LLM extraction fail schema → rơi vào review queue, không bao giờ crash pipeline.
6.2 Flow F — Price-change detection
Cron re-crawl campaign đang OPEN → extract giá mới
→ lệch > ngưỡng (vd 5%) → Telegram alert + banner warning storefront
→ KHÔNG tự sửa giá campaign đang chạy — quyết định thuộc creator/admin7 · TIỀN: CHỐT GIÁ · SETTLEMENT · HOÀN TIỀN
Quyết định business: final ≥ deposit là tiên quyết ⇒ settlement chỉ thu thêm, không hoàn chênh lệch. Nhờ đó đã XÓA khỏi scope: adjustment 2 chiều, refund-difference jobs, trạng thái UI "hoàn tiền dương". Chỉ thêm 1 worker + 1 bảng.
7.1 Order State Machine V2 (vòng đời 2 pha thanh toán)
┌──────────────┐
đặt chỗ 15' ──────▶│ RESERVED │───── hết hạn ─────▶ EXPIRED
└──────┬───────┘
webhook cọc (idempotent)
▼
┌──────────────┐
│ DEPOSIT_HELD │◀──── đơn này mới được ĐẾM vào MOQ/tier
└──────┬───────┘
campaign CLOSED → khoá tier (atomic)
▼
┌──────────────────┐ ◀── SettlementBatch sinh bill
│ SETTLEMENT_DUE │ amount_due ≥ 0 (DB CHECK)
└───┬──────────┬───┘
trả phần còn │ │ quá deadline (+grace 10')
lại đúng hạn ▼ ▼
┌────────────┐ ┌─────────────────────┐
│ CONFIRMED │ │ CANCELLED_NO_SETTLE │
└─────┬──────┘ └──────────┬──────────┘
▼ ▼ hoàn CỌC
FULFILLING ──▶ SHIPPED REFUNDING ──▶ REFUNDED
(RefundEngine resumable)7.2 Luồng tiền 2 pha
PHA 1 · CỌC (lúc đang OPEN) PHA 2 · CHỐT (lúc CLOSE)
┌─────────────────────────────┐ ┌──────────────────────────────┐
│ deposit = pct% × GIÁ TỐT NHẤT│ │ tổng qty thực tế → tra tier │
│ (tier maxQty — rẻ nhất) │ │ final_unit_price = tier.price │
│ │ │ │
│ Neo giá rẻ nhất ⇒ mọi kịch │ ──────▶ │ amount_due = │
│ bản final ∈ [rẻ…đắt] ⇒ │ │ qty × final │
│ final ≥ cọc LUÔN đúng ✅ │ │ + shipping_fee │
└─────────────────────────────┘ │ − deposit_paid = CÒN LẠI ≥ 0 │
└──────────────┬───────────────┘
▼
user chuyển khoản/QR
▼
✅ CONFIRMED → vận chuyểnCopy UI storefront gợi ý:"Cọc {X}đ giữ slot · Giá chốt sau khi đóng nhóm, dao động {P_đắt}đ–{P_rẻ}đ/sản phẩm · Cọc được trừ thẳng vào giá chốt."
7.3 Khoá tier lúc CLOSE (atomic)
sql
BEGIN;
SELECT qty FROM orders
WHERE campaign_id = :cid AND status = 'DEPOSIT_HELD'
FOR UPDATE; -- đóng băng tổng qty
-- app: SUM(qty) → tra bảng tiers → tier_id
UPDATE campaigns SET status='CLOSED', locked_tier_id = :tier_id WHERE id = :cid;
INSERT INTO outbox (event_type, payload) VALUES ('campaign.tier_locked', ...);
COMMIT;7.4 Flow G — Settlement Batch (sinh bill hàng loạt, resumable)
outbox 'campaign.tier_locked' ──▶ enqueue SettlementBatch{campaign_id, cursor:null}
│
▼
Worker loop (chunk 100):
INSERT INTO settlements (...) SELECT ... FROM orders WHERE status='DEPOSIT_HELD'
ON CONFLICT (order_id) DO NOTHING; ← chạy lại 100 lần vẫn an toàn
UPDATE orders SET status='SETTLEMENT_DUE';
INSERT outbox('settlement.due') ──▶ Notifier: TG/email cho user
"Bạn cần chuyển {amount_due}đ trước {deadline} · QR kèm message {order_code}"
checkpoint cursor vào PG.refund_batches-style tracking7.5 Flow H — Settlement window & ca biên
Deadline worker (repeatable scan mỗi phút):
sql
SELECT id FROM settlements
WHERE status='DUE' AND deadline + INTERVAL '10 minutes' < now()
FOR UPDATE SKIP LOCKED LIMIT 100; -- nhiều instance song song vẫn ổn
-- → tx: settlement='CANCELLED' + order='CANCELLED_NO_SETTLE'
-- → outbox('order.cancelled') → RefundEngine hoàn cọcGuard bổ sung trong ProcessPayment:
| Ca | Xử lý |
|---|---|
| Tiền phần-còn-lại đến đúng giây deadline | Row lock + recheck: now() ≤ deadline + grace thì chấp nhận — ai đến trước thắng |
| Webhook cọc đến SAU khi campaign CLOSED | Không cộng vào MOQ → auto hoàn toàn bộ + notify "nhóm đã đóng" |
| Webhook cọc đến khi campaign CANCELLED | Auto hoàn + notify |
7.6 Flow E — RefundEngine (resumable, chống double-refund)
outbox 'campaign.cancelled' hoặc 'order.cancelled'
→ enqueue RefundBatch{scope, cursor:null}
Worker loop (chunk 100):
SELECT id FROM orders WHERE <scope> AND refund_status IS NULL ORDER BY id LIMIT 100;
WITH EACH:
INSERT refunds (order_id, idem_key='refund:'||order_id, ...)
ON CONFLICT (idem_key) DO NOTHING; ← chống double-refund tuyệt đối
→ gọi provider refund (kèm idem key)
→ mark refund_status='INITIATED'
checkpoint cursor vào PG
crash bất kỳ đâu → retry → resume từ cursor, không mất/không lặp7.7 Reconciliation cron (hằng ngày)
- Đối soát
payments↔ statement của provider → lệch → Telegram alert. - Invariant check:
SUM(payments WHERE kind='DEPOSIT' AND ok) == SUM(orders.deposit_paid).
8 · DATA MODEL ĐẦY ĐỦ
8.1 Commerce lõi
sql
CREATE TYPE campaign_status AS ENUM
('DRAFT','PENDING_REVIEW','OPEN','CLOSED','FULFILLING','DONE','CANCELLED');
CREATE TABLE campaigns (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
title TEXT NOT NULL,
status campaign_status NOT NULL DEFAULT 'DRAFT',
deposit_mode TEXT NOT NULL DEFAULT 'PERCENT_OF_CHEAPEST',
deposit_per_unit NUMERIC(12,0) NOT NULL,
settlement_window_h INT NOT NULL DEFAULT 72,
grace_minutes INT NOT NULL DEFAULT 10,
shipping_fee NUMERIC(12,0) NOT NULL DEFAULT 0,
closes_at TIMESTAMPTZ,
locked_tier_id BIGINT,
capacity INT NOT NULL,
taken INT NOT NULL DEFAULT 0,
version INT NOT NULL DEFAULT 0, -- optimistic lock
CHECK (taken >= 0 AND taken <= capacity),
CHECK (deposit_per_unit >= 0)
);
CREATE TABLE campaign_tiers (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
campaign_id BIGINT NOT NULL REFERENCES campaigns(id),
qty_threshold INT NOT NULL, -- đạt N đơn mới hưởng giá này
unit_price NUMERIC(12,0) NOT NULL,
UNIQUE (campaign_id, qty_threshold)
);
CREATE TYPE order_status AS ENUM
('RESERVED','EXPIRED','DEPOSIT_HELD','SETTLEMENT_DUE',
'CONFIRMED','FULFILLING','SHIPPED',
'CANCELLED','CANCELLED_NO_SETTLE','REFUNDING','REFUNDED');
CREATE TABLE orders (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
code TEXT NOT NULL UNIQUE, -- mã chuyển khoản "GB-XXXX"
user_id BIGINT NOT NULL,
campaign_id BIGINT NOT NULL REFERENCES campaigns(id),
qty INT NOT NULL CHECK (qty > 0),
status order_status NOT NULL DEFAULT 'RESERVED',
expires_at TIMESTAMPTZ, -- hạn giữ chỗ 15'
deposit_due NUMERIC(12,0) NOT NULL,
deposit_paid NUMERIC(12,0) NOT NULL DEFAULT 0,
settlement_paid NUMERIC(12,0) NOT NULL DEFAULT 0,
refund_status TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX idx_orders_campaign_status ON orders (campaign_id, status);
-- LƯỚI AN TOÀN: dù app bug cũng không ra tiền âm
ALTER TABLE orders ADD CONSTRAINT chk_money_non_negative
CHECK (deposit_paid >= 0 AND settlement_paid >= 0);8.2 Tiền & webhook
sql
CREATE TABLE payments (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
order_id BIGINT NOT NULL REFERENCES orders(id),
kind TEXT NOT NULL CHECK (kind IN ('DEPOSIT','REMAINDER','REFUND')),
amount NUMERIC(12,0) NOT NULL CHECK (amount >= 0),
provider TEXT NOT NULL,
provider_ref TEXT,
status TEXT NOT NULL DEFAULT 'PENDING'
CHECK (status IN ('PENDING','SUCCEEDED','FAILED')),
idempotency_key TEXT NOT NULL UNIQUE,
raw JSONB,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE webhook_events (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
provider TEXT NOT NULL,
idempotency_key TEXT NOT NULL UNIQUE,
payload JSONB NOT NULL,
processed_at TIMESTAMPTZ
);
CREATE TABLE settlements (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
order_id BIGINT NOT NULL UNIQUE REFERENCES orders(id),
final_unit_price NUMERIC(12,0) NOT NULL, -- giá tier đã khoá
qty INT NOT NULL,
goods_total NUMERIC(12,0) NOT NULL, -- final × qty
shipping_fee NUMERIC(12,0) NOT NULL DEFAULT 0,
deposit_credit NUMERIC(12,0) NOT NULL,
amount_due NUMERIC(12,0) GENERATED ALWAYS AS
(goods_total + shipping_fee - deposit_credit) STORED,
deadline TIMESTAMPTZ NOT NULL,
status TEXT NOT NULL DEFAULT 'DUE'
CHECK (status IN ('DUE','PAID','CANCELLED'))
);
ALTER TABLE settlements ADD CONSTRAINT chk_due_non_negative
CHECK (goods_total + shipping_fee - deposit_credit >= 0);
CREATE INDEX idx_settlements_scan ON settlements (status, deadline);
CREATE TABLE refunds (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
order_id BIGINT NOT NULL REFERENCES orders(id),
amount NUMERIC(12,0) NOT NULL,
idem_key TEXT NOT NULL UNIQUE, -- 'refund:'||order_id
status TEXT NOT NULL DEFAULT 'INITIATED',
provider_ref TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE refund_batches (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
scope JSONB NOT NULL, -- {campaign_id} hoặc {order_id}
cursor BIGINT,
status TEXT NOT NULL DEFAULT 'RUNNING'
CHECK (status IN ('RUNNING','COMPLETE','FAILED')),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);8.3 Outbox & audit
sql
CREATE TABLE outbox (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
aggregate_type TEXT NOT NULL,
aggregate_id BIGINT NOT NULL,
event_type TEXT NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
published_at TIMESTAMPTZ
);
CREATE INDEX idx_outbox_unpublished ON outbox (id) WHERE published_at IS NULL;
CREATE TABLE audit_log (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
actor TEXT NOT NULL, -- user/admin/system
entity TEXT NOT NULL,
entity_id BIGINT NOT NULL,
action TEXT NOT NULL,
diff JSONB,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);8.4 Crawler
sql
CREATE TABLE crawl_sources (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
name TEXT NOT NULL,
type TEXT NOT NULL CHECK (type IN ('SITE','CHAT')),
target TEXT NOT NULL, -- URL / channel
schedule TEXT NOT NULL, -- cron expr
rate_limit REAL NOT NULL DEFAULT 0.5, -- req/s theo domain
trust_level INT NOT NULL DEFAULT 0, -- whitelisted auto-publish?
enabled BOOLEAN NOT NULL DEFAULT TRUE
);
CREATE TABLE raw_snapshots (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
source_id BIGINT NOT NULL REFERENCES crawl_sources(id),
s3_key TEXT NOT NULL, -- html.gz / screenshot.png
content_hash TEXT NOT NULL,
fetched_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TYPE draft_status AS ENUM
('EXTRACTING','PENDING_REVIEW','APPROVED','REJECTED');
CREATE TABLE campaign_drafts (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
snapshot_id BIGINT NOT NULL REFERENCES raw_snapshots(id),
extraction JSONB NOT NULL, -- output LLM
confidence REAL,
status draft_status NOT NULL DEFAULT 'PENDING_REVIEW',
campaign_id BIGINT REFERENCES campaigns(id), -- sau khi approve & publish
reviewed_by TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE price_checks (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
campaign_id BIGINT NOT NULL REFERENCES campaigns(id),
detected_price NUMERIC(12,0) NOT NULL,
delta_pct REAL NOT NULL,
alerted_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);9 · BUSINESS KNOBS
| Tham số | Default | Ghi chú |
|---|---|---|
deposit_mode | PERCENT_OF_CHEAPEST = 30% | Bảo đảm invariant final ≥ cọc tự nhiên |
reservation_ttl | 15 phút | Giữ slot chưa cọc |
settlement_window | 72h | Cửa sổ trả phần còn lại |
grace_period | 10 phút | Chống dispute ở biên deadline |
unpaid_action | CANCEL + FULL_REFUND_DEPOSIT | V1 không trừ phí |
shipping_fee | Flat / theo vùng | Cộng thẳng vào amount_due |
price_alert_threshold | 5% | Ngưỡng cảnh báo giá nguồn đổi |
auto_publish_confidence | 0.95 (mặc định OFF) | Chỉ bật với source whitelist |
10 · ROADMAP TRIỂN KHAI
Chi tiết đầy đủ: xem
architecture-roadmap.md
| Phase | Tên | Nội dung | Tiêu chí xong (acceptance) |
|---|---|---|---|
| 0 | First Order | Go monolith: Auth, Campaign, Order, Deposit Expiry (in-proc), VietQR webhook, Settlement cơ bản | Bán được đơn đầu tiên cho shop mình. Luồng A→Z chạy an toàn |
| 1 | Operational | Asynq thay in-proc cron, Telegram Bot notify, Redis, webhook async, structured logging | Hệ thống tự chạy, không cần ngồi canh |
| 2 | Realtime | SSE slot counter, Rate Limiting, Redis cache, Health check | 50+ concurrent OK, realtime đếm slot |
| 3 | Settlement | Tier locking, SettlementBatch, Deadline Scanner, RefundEngine, Reconciliation, Outbox | Đóng nhóm tự động, thu tiền tự động, hoàn cọc tự động |
| 4 | Scale | Tách Order Engine (gRPC), Worker riêng, PgBouncer, Casbin, OTel | 1000+ concurrent, trace xuyên suốt |
| 5 | Platform | Multi-tenant (RLS), Workspace RBAC, Crawler Pipeline, Billing | Nhiều shop cùng chạy trên 1 hệ thống |
11 · CHẠY LOCAL: DOCKER-COMPOSE & REPO LAYOUT
Compose services
web(gateway) engine(go) workers crawler
postgres pgbouncer valkey minio
otel-collector jaeger prometheus grafana ollama⚠️ PgBouncer chế độ
transaction+ pgx: setQueryExecModeCacheStatementđể tránh lỗi prepared statements.
Repo layout (monorepo)
groupbuy/
├── apps/
│ ├── web/ # Next.js
│ ├── gateway/ # NestJS (+ orchestrator)
│ ├── engine/ # Go
│ ├── workers/ # BullMQ fleet
│ └── crawler/ # Playwright + LLM
├── packages/
│ ├── proto/ # .proto — hợp đồng gRPC dùng chung
│ ├── schemas/ # Zod schemas share FE/BE/crawler
│ └── otel/ # tracing bootstrap dùng lại mọi service
└── infra/
├── docker-compose.yml
└── otel/collector.yaml12 · LỘ TRÌNH HỌC
| Tuần | Mục tiêu | "Thấy được đồ vật" |
|---|---|---|
| 1 | docker-compose PG + Valkey + Jaeger + Grafana; Gateway skeleton + auto-instrumentation | Mở localhost:16686 — trace đầu tiên 🎉 |
| 2 | Viết .proto + Go engine + otelpgx | Trace chảy từ NestJS sang Go |
| 3 | Checkout E2E: ReserveSlot → webhook cọc → outbox → SSE counter | Slot counter realtime trên 2 browser |
| 4 | Workers + BullMQ trace injection | Trace nối tiếp qua job delay 15 phút |
| 5 | Close → SettlementBatch → deadline → refund | Demo full vòng đời tiền 2 pha |
| 6+ | Crawler + Ollama extraction + Draft Review UI | Campaign đầu tiên sinh từ crawl |
Hết tài liệu — ARCHITECTURE.md v3.0