Skip to content

🏗️ 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

Hạng mụcQuyết định
Sản phẩmNề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 doanhfinal_price ≥ deposit luôn đúng ⇒ hệ thống chỉ thu thêm, không bao giờ hoàn chênh lệch
Hot pathGo Order Engine — chống oversell bằng atomic SQL trên PostgreSQL
ConsistencyPostgreSQL là source of truth duy nhất; mọi state transition đi qua Outbox Pattern
TiềnWebhook idempotent · settlement batch resumable · hoàn cọc tự động khi user bỏ qua hạn
Tự động hoá nguồn hàngCrawler (Playwright) → LLM extract → Dedup → Draft PENDING_REVIEW → Approve → Publish
Quan sátOpenTelemetry 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

  1. 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.
  2. Outbox Pattern: mỗi transition ghi state + event trong cùng 1 transaction; side effects chỉ dispatch qua event/job.
  3. Webhook idempotency-first: verify HMAC → insert với unique key → trả 200 ngay → xử lý async.
  4. Chỉ đơn DEPOSIT_HELD mới đếm vào MOQ/tier ⇒ tổngqty ổn định tuyệt đối tại thời điểm close.
  5. deposit_per_unit ≤ MIN(tier.unit_price) ⇒ mọi kịch bản chốt đều final ≥ cọc.
  6. Settlement 1 chiều (chỉ thu thêm) — xoá khỏi scope nhánh hoàn-chênh-lệch.
  7. Quá hạn settlement → tự huỷ + hoàn cọc qua RefundEngine (pattern resumable dùng lại).
  8. Crawler không bao giờ auto-publish 100%: luôn qua PENDING_REVIEW; auto-publish chỉ theo confidence threshold + global kill-switch.
  9. Observability từ ngày 1, chạy ngang mọi layer (cross-cutting).
  10. 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ồngGiao thức
Client → Next.jsHTTPS
Storefront → GatewayREST JSON + SSE stream
Admin → GatewayGraphQL (Mercurius)
Gateway ↔ Go EnginegRPC + Protobuf
Mọi service → hàng đợiBullMQ (Valkey/Redis)
Outbox relay → OrchestratorPoll PG 1s (sau nâng cấp Debezium)
Cron scheduler → CrawlerBullMQ repeatable jobs
Ra thế giới ngoàiHTTPS API (Stripe/VNPay, Telegram, LLM…)
Services → DBTCP (pgx / ioredis / S3 API)
Services → CollectorOTLP 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ầnChọnGiấy phép / Free tier
Frontend hostCloudflare PagesFree unlimited BW
API frameworkNestJS + Fastify + MercuriusMIT
Authjose + RedisMIT
Hot pathGo + grpc + pgx + sqlcBSD/MIT
QueueBullMQ + Valkey (hoặc Redis 7)MIT / BSD
DatabasePostgres 16 + PgBouncer (Docker hoặc Neon free)OSS
Object storageMinIO local / Cloudflare R2 (10GB free)AGPL / free tier
CrawlerPlaywright + CheerioApache-2.0 / MIT
LLMOllama + qwen2.5:7b (local) / Groq free tierMIT / free
Payment testStripe test mode / VNPay sandboxFree
NotifyTelegram Bot API + Resend (3k email/tháng)Free
TracingOTel SDK + Collector + JaegerApache-2.0/CNCF
Metrics/DashboardPrometheus + GrafanaApache-2.0/AGPL
Logspino (+ Loki nếu muốn tập trung)MIT/AGPL
Error trackingGlitchTip 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 slots

4.2 Instrumentation từng service

ServiceSDK / LibGhi chú
NestJS Gateway@opentelemetry/sdk-node + auto-instrumentations-nodeAuto bắt http, graphql, ioredis, pg
Worker FleetCùng SDK + inject thủ công qua job.data.__otelBullMQ không tự propagate
Go Enginego.opentelemetry.io/otel + otelgrpc + otelpgxInterceptor gRPC 2 phía
FrontendBrowser tự gửi header traceparentKhông cần SDK nặng
Collectorotelcol-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 trace

4.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 RETURNINGatomic single-statement — nhanh hơn SELECT FOR UPDATE vì 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êng

Nguyên tắc an toàn:

  • Giai đoạn đầu: human review 100%. Auto-publish chỉ bật khi confidence ≥ 0.95 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/admin

7 · 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ển

Copy 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 tracking

7.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ọc

Guard bổ sung trong ProcessPayment:

CaXử lý
Tiền phần-còn-lại đến đúng giây deadlineRow lock + recheck: now() ≤ deadline + grace thì chấp nhận — ai đến trước thắng
Webhook cọc đến SAU khi campaign CLOSEDKhông cộng vào MOQ → auto hoàn toàn bộ + notify "nhóm đã đóng"
Webhook cọc đến khi campaign CANCELLEDAuto 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ặp

7.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ốDefaultGhi chú
deposit_modePERCENT_OF_CHEAPEST = 30%Bảo đảm invariant final ≥ cọc tự nhiên
reservation_ttl15 phútGiữ slot chưa cọc
settlement_window72hCửa sổ trả phần còn lại
grace_period10 phútChống dispute ở biên deadline
unpaid_actionCANCEL + FULL_REFUND_DEPOSITV1 không trừ phí
shipping_feeFlat / theo vùngCộng thẳng vào amount_due
price_alert_threshold5%Ngưỡng cảnh báo giá nguồn đổi
auto_publish_confidence0.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

PhaseTênNội dungTiêu chí xong (acceptance)
0First OrderGo monolith: Auth, Campaign, Order, Deposit Expiry (in-proc), VietQR webhook, Settlement cơ bảnBán được đơn đầu tiên cho shop mình. Luồng A→Z chạy an toàn
1OperationalAsynq thay in-proc cron, Telegram Bot notify, Redis, webhook async, structured loggingHệ thống tự chạy, không cần ngồi canh
2RealtimeSSE slot counter, Rate Limiting, Redis cache, Health check50+ concurrent OK, realtime đếm slot
3SettlementTier locking, SettlementBatch, Deadline Scanner, RefundEngine, Reconciliation, OutboxĐóng nhóm tự động, thu tiền tự động, hoàn cọc tự động
4ScaleTách Order Engine (gRPC), Worker riêng, PgBouncer, Casbin, OTel1000+ concurrent, trace xuyên suốt
5PlatformMulti-tenant (RLS), Workspace RBAC, Crawler Pipeline, BillingNhiề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: set QueryExecModeCacheStatement để 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.yaml

12 · LỘ TRÌNH HỌC

TuầnMục tiêu"Thấy được đồ vật"
1docker-compose PG + Valkey + Jaeger + Grafana; Gateway skeleton + auto-instrumentationMở localhost:16686trace đầu tiên 🎉
2Viết .proto + Go engine + otelpgxTrace chảy từ NestJS sang Go
3Checkout E2E: ReserveSlot → webhook cọc → outbox → SSE counterSlot counter realtime trên 2 browser
4Workers + BullMQ trace injectionTrace nối tiếp qua job delay 15 phút
5Close → SettlementBatch → deadline → refundDemo full vòng đời tiền 2 pha
6+Crawler + Ollama extraction + Draft Review UICampaign đầu tiên sinh từ crawl

Hết tài liệu — ARCHITECTURE.md v3.0