#6 Bước 4 — Dựng module `identity` (Quản lý User & Auth) Phần 2
Bảng ánh xạ, để tiện đối chiếu:
| Lỗi | HTTP | Vì sao |
|---|---|---|
ErrInvalidInput |
422 Unprocessable Entity |
JSON đúng cú pháp, nhưng nội dung vi phạm quy tắc |
ErrEmailTaken / ErrUsernameTaken |
409 Conflict |
Xung đột với trạng thái hiện có của hệ thống |
ErrInvalidCredentials |
401 Unauthorized |
Chưa chứng minh được danh tính |
ErrUserNotFound |
404 Not Found |
|
| JSON hỏng | 400 Bad Request |
Không phân tích nổi request |
| Còn lại | 500 |
Lỗi của chúng ta, không phải của client |
400hay422?400nghĩa là "tôi không đọc nổi request này".422nghĩa là "tôi đọc được, hiểu được, nhưng nội dung sai". Phân biệt được hai cái giúp client biết nên sửa cách gửi hay sửa dữ liệu người dùng nhập.
14.4 Router
Tạo internal/modules/identity/transport/http/router.go:
package http
import "github.com/go-chi/chi/v5"
// Mount gắn toàn bộ route của module vào router cha.
//
// Đọc hàm này là biết ngay module phơi ra những gì và cái nào cần đăng nhập.
func (h *Handler) Mount(r chi.Router) {
// Công khai
r.Post("/auth/register", h.handleRegister)
r.Post("/auth/login", h.handleLogin)
// Cần bearer token. Group tạo một nhánh middleware riêng, không ảnh
// hưởng tới các route đã đăng ký ở trên.
r.Group(func(protected chi.Router) {
protected.Use(middlewareAuth(h))
protected.Get("/me", h.handleMe)
})
}
Thêm vào cuối handler.go:
func middlewareAuth(h *Handler) func(nextHandler http.Handler) http.Handler {
return middleware.Auth(h.verifier)
}
Hàm bọc một dòng này chỉ để
router.gokhông phải import thêm package. Nếu thấy thừa, gọi thẳngmiddleware.Auth(h.verifier)trongMountcũng đúng.
15. Lắp ráp: module.go và cmd/api
15.1 module.go — điểm lắp ráp duy nhất
Tạo internal/modules/identity/module.go:
// Package identity là điểm lắp ráp của module quản lý người dùng.
//
// Đây là package DUY NHẤT bên ngoài được phép import. Mọi thứ trong
// domain/, service/, repository/, transport/ là chuyện nội bộ.
package identity
import (
"log/slog"
"github.com/go-chi/chi/v5"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/yourname/community/internal/modules/identity/repository"
"github.com/yourname/community/internal/modules/identity/service"
identityhttp "github.com/yourname/community/internal/modules/identity/transport/http"
"github.com/yourname/community/internal/platform/httpx/middleware"
"github.com/yourname/community/internal/platform/outbox"
"github.com/yourname/community/internal/platform/postgres"
)
type Module struct {
api *identityhttp.Handler
}
func New(
pool *pgxpool.Pool,
txm *postgres.TxManager,
verifier middleware.Verifier,
tokens service.TokenIssuer,
log *slog.Logger,
) *Module {
repo := repository.NewPostgres(pool)
svc := service.New(txm, repo, outbox.NewWriter(), tokens, log)
return &Module{
api: identityhttp.NewHandler(svc, verifier, log),
}
}
func (m *Module) RegisterHTTP(r chi.Router) {
m.api.Mount(r)
}
Trong thực tế verifier và tokens là cùng một *token.Issuer; hai tham số riêng vì đó là hai vai trò khác nhau (cấp token / kiểm token), và ngày bạn tách khoá ký khỏi khoá kiểm (chuyển sang RS256) thì chữ ký hàm không phải đổi.
Chưa có RegisterEvents vì module này chưa nghe event nào. Nó sẽ xuất hiện khi có việc thật để nghe.
15.2 cmd/api/main.go
package main
import (
"context"
"errors"
"fmt"
"log/slog"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/go-chi/chi/v5"
"github.com/yourname/community/internal/modules/identity"
"github.com/yourname/community/internal/platform/config"
"github.com/yourname/community/internal/platform/httpx"
"github.com/yourname/community/internal/platform/httpx/middleware"
"github.com/yourname/community/internal/platform/logger"
"github.com/yourname/community/internal/platform/postgres"
"github.com/yourname/community/internal/platform/token"
)
func main() {
if err := run(); err != nil {
slog.Error("api dừng vì lỗi", "err", err)
os.Exit(1)
}
}
func run() error {
cfg, err := config.Load()
if err != nil {
return err
}
log := logger.New(cfg.LogLevel)
slog.SetDefault(log)
// ctx bị huỷ khi nhận Ctrl+C hoặc SIGTERM (docker stop, kubectl delete).
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
pool, err := postgres.Connect(ctx, postgres.DefaultConfig(cfg.DatabaseURL, "community-api"))
if err != nil {
return err
}
defer pool.Close()
tokens, err := token.NewIssuer(cfg.JWTSecret, cfg.JWTTTL, "community")
if err != nil {
return err
}
txm := postgres.NewTxManager(pool)
r := chi.NewRouter()
r.Use(middleware.RequestID)
r.Use(middleware.Recover(log))
r.Get("/healthz", func(w http.ResponseWriter, _ *http.Request) {
httpx.JSON(w, http.StatusOK, map[string]string{"status": "ok"})
})
// Danh sách module. Thêm module mới = thêm một dòng ở đây.
identity.New(pool, txm, tokens, tokens, log).RegisterHTTP(r)
srv := &http.Server{
Addr: fmt.Sprintf(":%d", cfg.HTTPPort),
Handler: r,
// ReadHeaderTimeout chặn Slowloris: kết nối gửi header nhỏ giọt
// vài byte mỗi phút để giữ chỗ mãi mãi. Không có nó, vài trăm
// kết nối như vậy là đủ làm server ngừng nhận request mới.
ReadHeaderTimeout: 5 * time.Second,
ReadTimeout: 15 * time.Second,
WriteTimeout: 15 * time.Second,
IdleTimeout: 60 * time.Second,
}
serverErr := make(chan error, 1)
go func() {
log.Info("api đang lắng nghe", "addr", srv.Addr)
if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
serverErr <- err
}
close(serverErr)
}()
select {
case err := <-serverErr:
return err
case <-ctx.Done():
log.Info("nhận tín hiệu dừng, đang đóng kết nối")
}
// Dùng context MỚI: ctx đã bị huỷ, dùng lại nó thì Shutdown ngắt ngay
// lập tức và giết cả những request đang xử lý dở.
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
return srv.Shutdown(shutdownCtx)
}
identity.New(pool, txm, tokens, tokens, log)—tokensxuất hiện hai lần, đúng như §15.1 giải thích: một lần với vai tròVerifier, một lần với vai tròTokenIssuer.
Vì sao
mainchỉ có ba dòng cònrunlàm hết? Vìos.Exitkhông chạy các hàmdefer. Nếumainvừa gọios.Exitvừadefer pool.Close(), kết nối DB không bao giờ được đóng. Tách ra:runtrả lỗi và đóng gọn mọi thứ,mainchỉ quyết định mã thoát.
15.3 Cập nhật Makefile
api: ## Chạy HTTP API
go run ./cmd/api
test: ## Chạy toàn bộ test
go test ./internal/... -race -count=1
16. Chạy thử đầu-cuối
16.1 Khởi động
docker compose -f deployments/docker-compose.yml up -d
. .\dev.ps1
Import-DotEnv
go run ./cmd/api
{"time":"2026-08-01T09:12:03Z","level":"INFO","msg":"api đang lắng nghe","addr":":8000"}
16.2 Đăng ký
Mở terminal thứ hai:
$body = @{
name = 'Nguyễn Huy Hoàng'
username = 'hhoang'
email = 'hhoang@example.com'
password = 'matkhau-du-dai-123'
} | ConvertTo-Json
$res = Invoke-RestMethod -Uri http://localhost:8000/auth/register `
-Method Post `
-Body ([Text.Encoding]::UTF8.GetBytes($body)) `
-ContentType 'application/json; charset=utf-8'
$res | ConvertTo-Json -Depth 3
⚠️
[Text.Encoding]::UTF8.GetBytes($body)không phải thừa. Windows PowerShell 5.1 gửi chuỗi theo bảng mã ISO-8859-1, nên"Nguyễn Huy Hoàng"sẽ tới server thành"Nguyá»…n Huy Hoà ng". Chuyển sang mảng byte UTF-8 là cách chắc chắn nhất. (PowerShell 7 đã sửa mặc định này.)
{
"user": {
"id": "5c8f1b2e-...",
"name": "Nguyễn Huy Hoàng",
"username": "hhoang",
"email": "hhoang@example.com",
"avatar_url": null,
"reputation": 0,
"created_at": "2026-08-01T09:12:31Z"
},
"token": "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9...",
"expires_at": "2026-08-02T09:12:31Z"
}
Không có password hay password_hash trong response — đó là UserResponse ở §14.2 làm việc.
16.3 Kiểm tra event đã vào outbox
Đây là phần đáng xem nhất của cả bước:
docker compose -f deployments/docker-compose.yml exec -T postgres `
psql -U app -d community -c `
"SELECT id, event_type, aggregate_id, published_at, attempts FROM outbox ORDER BY id DESC LIMIT 1;"
id | event_type | aggregate_id | published_at | attempts
----+------------------------------+--------------------------------------+--------------+----------
1 | identity.user.registered.v1 | 5c8f1b2e-... | | 0
published_at rỗng. Event đang nằm chờ, chưa ai gửi nó đi đâu cả — vì relay là bước 5. Đúng như thiết kế: tầng ghi không nói chuyện với Kafka, nó chỉ ghi vào một bảng.
Xem nội dung payload:
docker compose -f deployments/docker-compose.yml exec -T postgres `
psql -U app -d community -c "SELECT jsonb_pretty(payload) FROM outbox ORDER BY id DESC LIMIT 1;"
{
"name": "Nguyễn Huy Hoàng",
"email": "hhoang@example.com",
"user_id": "5c8f1b2e-...",
"username": "hhoang",
"registered_at": "2026-08-01T09:12:31.482913Z"
}
Kiểm tra correlation_id — nó phải khớp với header X-Request-ID của chính request đăng ký vừa rồi:
docker compose -f deployments/docker-compose.yml exec -T postgres `
psql -U app -d community -c "SELECT correlation_id FROM outbox ORDER BY id DESC LIMIT 1;"
Đây là sợi chỉ sẽ nối HTTP request → outbox → Kafka → consumer. Khi có sự cố ở bước 6, một chuỗi này là đủ để lần lại toàn bộ hành trình.
16.4 Các nhánh lỗi
# Trùng email → 409
Invoke-RestMethod -Uri http://localhost:8000/auth/register -Method Post `
-Body ([Text.Encoding]::UTF8.GetBytes($body)) -ContentType 'application/json; charset=utf-8'
{"error":{"code":"email_taken","message":"email đã được sử dụng"},"request_id":"..."}
Kiểm tra ngay: bảng outbox vẫn chỉ có một dòng. Lần đăng ký thất bại không để lại event nào — transaction đã rollback cả hai bảng.
# Mật khẩu ngắn → 422
$bad = @{ name='A'; username='abc'; email='a@b.com'; password='123' } | ConvertTo-Json
Invoke-RestMethod -Uri http://localhost:8000/auth/register -Method Post `
-Body ([Text.Encoding]::UTF8.GetBytes($bad)) -ContentType 'application/json; charset=utf-8'
{"error":{"code":"invalid_input","message":"dữ liệu không hợp lệ: mật khẩu phải từ 8 đến 72 byte"}}
16.5 Đăng nhập và gọi /me
$login = @{ email='HHoang@Example.COM'; password='matkhau-du-dai-123' } | ConvertTo-Json
$auth = Invoke-RestMethod -Uri http://localhost:8000/auth/login -Method Post `
-Body ([Text.Encoding]::UTF8.GetBytes($login)) -ContentType 'application/json; charset=utf-8'
Invoke-RestMethod -Uri http://localhost:8000/me `
-Headers @{ Authorization = "Bearer $($auth.token)" }
Chú ý email viết hoa lung tung vẫn đăng nhập được — NormalizeEmail đã xử lý đúng cái bẫy ở bước 3 §3.3.
Thử token hỏng:
Invoke-RestMethod -Uri http://localhost:8000/me -Headers @{ Authorization = "Bearer khong-phai-token" }
# 401 {"error":{"code":"unauthorized","message":"token không hợp lệ"}}
17. Kiểm chứng bằng test
17.1 password
Tạo internal/platform/password/password_test.go:
package password_test
import (
"strings"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/yourname/community/internal/platform/password"
)
func TestHashVaVerify(t *testing.T) {
h, err := password.Hash("matkhau-bi-mat")
require.NoError(t, err)
assert.NotContains(t, h, "matkhau", "hash không được chứa mật khẩu gốc")
assert.NoError(t, password.Verify(h, "matkhau-bi-mat"))
assert.ErrorIs(t, password.Verify(h, "sai-mat-khau"), password.ErrMismatch)
}
func TestHaiLanBamCungMatKhauChoHaiHashKhacNhau(t *testing.T) {
a, err := password.Hash("giong-nhau")
require.NoError(t, err)
b, err := password.Hash("giong-nhau")
require.NoError(t, err)
// bcrypt tự sinh salt ngẫu nhiên. Nhờ vậy hai người dùng đặt cùng
// mật khẩu vẫn có hash khác nhau, và bảng tra sẵn (rainbow table)
// trở nên vô dụng.
assert.NotEqual(t, a, b)
assert.NoError(t, password.Verify(a, "giong-nhau"))
assert.NoError(t, password.Verify(b, "giong-nhau"))
}
func TestMatKhauQuaDaiBiTuChoi(t *testing.T) {
_, err := password.Hash(strings.Repeat("a", 73))
assert.ErrorIs(t, err, password.ErrTooLong,
"phải từ chối rõ ràng thay vì cắt bớt trong im lặng")
}
17.2 token
Tạo internal/platform/token/token_test.go:
package token_test
import (
"strings"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/yourname/community/internal/platform/token"
)
const secret = "day-la-secret-du-32-ky-tu-cho-test!!"
func newIssuer(t *testing.T, ttl time.Duration) *token.Issuer {
t.Helper()
iss, err := token.NewIssuer(secret, ttl, "community")
require.NoError(t, err)
return iss
}
func TestCapRoiKiemLai(t *testing.T) {
iss := newIssuer(t, time.Hour)
raw, expiresAt, err := iss.Issue("user-123", time.Now())
require.NoError(t, err)
assert.WithinDuration(t, time.Now().Add(time.Hour), expiresAt, time.Minute)
sub, err := iss.Verify(raw)
require.NoError(t, err)
assert.Equal(t, "user-123", sub)
}
func TestTokenHetHanBiTuChoi(t *testing.T) {
iss := newIssuer(t, time.Hour)
// Cấp token với mốc thời gian của hai giờ trước — không phải chờ thật.
raw, _, err := iss.Issue("user-123", time.Now().Add(-2*time.Hour))
require.NoError(t, err)
_, err = iss.Verify(raw)
assert.ErrorIs(t, err, token.ErrInvalidToken)
}
func TestSecretKhacKhongKiemDuoc(t *testing.T) {
raw, _, err := newIssuer(t, time.Hour).Issue("user-123", time.Now())
require.NoError(t, err)
other, err := token.NewIssuer("mot-secret-hoan-toan-khac-du-32-byte", time.Hour, "community")
require.NoError(t, err)
_, err = other.Verify(raw)
assert.ErrorIs(t, err, token.ErrInvalidToken)
}
// Test quan trọng nhất của file này — xem tài liệu §8.3.
func TestTokenAlgNoneBiTuChoi(t *testing.T) {
// header {"alg":"none","typ":"JWT"} + payload {"sub":"admin"} + chữ ký rỗng
forged := "eyJhbGciOiJub25lIiwidHlwIjoiSldUIn0." +
"eyJzdWIiOiJhZG1pbiIsImlzcyI6ImNvbW11bml0eSJ9."
_, err := newIssuer(t, time.Hour).Verify(forged)
assert.ErrorIs(t, err, token.ErrInvalidToken,
"alg=none PHẢI bị từ chối — nếu test này đỏ, bất kỳ ai cũng đăng nhập được thành bất kỳ ai")
}
func TestSuaNoiDungLamHongChuKy(t *testing.T) {
iss := newIssuer(t, time.Hour)
raw, _, err := iss.Issue("user-123", time.Now())
require.NoError(t, err)
parts := strings.Split(raw, ".")
require.Len(t, parts, 3)
// Đổi payload nhưng giữ nguyên chữ ký cũ
tampered := parts[0] + ".eyJzdWIiOiJhZG1pbiJ9." + parts[2]
_, err = iss.Verify(tampered)
assert.ErrorIs(t, err, token.ErrInvalidToken)
}
func TestSecretQuaNganBiTuChoiNgayTuDau(t *testing.T) {
_, err := token.NewIssuer("ngan-qua", time.Hour, "community")
assert.Error(t, err, "lỗi cấu hình phải lộ ra lúc khởi động, không phải lúc phục vụ")
}
17.3 Test service không cần database
Tạo internal/modules/identity/service/service_test.go:
package service_test
import (
"context"
"testing"
"time"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/yourname/community/internal/contracts"
v1 "github.com/yourname/community/internal/contracts/v1"
"github.com/yourname/community/internal/modules/identity/domain"
"github.com/yourname/community/internal/modules/identity/service"
"github.com/yourname/community/internal/platform/logger"
)
// ─── Các bản giả ─────────────────────────────────────────────────────
// Chúng ngắn vì interface ở §13.1 chỉ khai đúng thứ service cần.
type fakeTx struct{ fail error }
func (f fakeTx) Do(_ context.Context, fn func(pgx.Tx) error) error {
if f.fail != nil {
return f.fail
}
return fn(nil) // repo và writer giả đều bỏ qua tham số tx
}
type fakeRepo struct {
byEmail map[string]domain.User
insertErr error
}
func (r *fakeRepo) Insert(_ context.Context, _ pgx.Tx, u domain.User) (domain.User, error) {
if r.insertErr != nil {
return domain.User{}, r.insertErr
}
u.CreatedAt = time.Now().UTC()
r.byEmail[u.Email] = u
return u, nil
}
func (r *fakeRepo) FindByEmail(_ context.Context, email string) (domain.User, error) {
u, ok := r.byEmail[email]
if !ok {
return domain.User{}, domain.ErrUserNotFound
}
return u, nil
}
func (r *fakeRepo) FindByID(context.Context, uuid.UUID) (domain.User, error) {
return domain.User{}, domain.ErrUserNotFound
}
type fakeWriter struct{ written []contracts.Envelope }
func (w *fakeWriter) Write(_ context.Context, _ pgx.Tx, _ string, e contracts.Envelope) error {
w.written = append(w.written, e)
return nil
}
type fakeTokens struct{}
func (fakeTokens) Issue(subject string, now time.Time) (string, time.Time, error) {
return "token-cho-" + subject, now.Add(time.Hour), nil
}
func newService(t *testing.T) (*service.Service, *fakeRepo, *fakeWriter) {
t.Helper()
repo := &fakeRepo{byEmail: map[string]domain.User{}}
writer := &fakeWriter{}
svc := service.New(fakeTx{}, repo, writer, fakeTokens{}, logger.New("error"))
return svc, repo, writer
}
func validRegistration() domain.Registration {
return domain.Registration{
Name: "Nguyễn Huy Hoàng",
Username: "HHoang", // chữ hoa — phải được chuẩn hoá
Email: " HHoang@Example.COM ", // hoa + khoảng trắng thừa
Password: "matkhau-du-dai-123",
}
}
// ─── Test ────────────────────────────────────────────────────────────
func TestDangKyChuanHoaEmailVaUsername(t *testing.T) {
svc, _, _ := newService(t)
res, err := svc.Register(context.Background(), validRegistration())
require.NoError(t, err)
assert.Equal(t, "hhoang@example.com", res.User.Email)
assert.Equal(t, "hhoang", res.User.Username)
assert.NotEmpty(t, res.Token)
}
// Nếu test này đỏ, module đã ghi user mà không phát event —
// và stats sẽ không bao giờ biết người dùng này tồn tại.
func TestDangKyPhatDungMotEventUserRegistered(t *testing.T) {
svc, _, writer := newService(t)
res, err := svc.Register(context.Background(), validRegistration())
require.NoError(t, err)
require.Len(t, writer.written, 1)
e := writer.written[0]
assert.Equal(t, v1.TypeUserRegistered, e.EventType)
assert.Equal(t, res.User.ID.String(), e.AggregateID)
assert.NoError(t, e.Validate())
payload, err := contracts.DecodePayload[v1.UserRegistered](e)
require.NoError(t, err)
assert.Equal(t, res.User.ID.String(), payload.UserID)
assert.Equal(t, "hhoang", payload.Username)
}
func TestDangKyLoiThiKhongPhatEvent(t *testing.T) {
svc, repo, writer := newService(t)
repo.insertErr = domain.ErrEmailTaken
_, err := svc.Register(context.Background(), validRegistration())
assert.ErrorIs(t, err, domain.ErrEmailTaken)
assert.Empty(t, writer.written, "insert thất bại thì không được có event nào")
}
func TestDangKyDuLieuSaiKhongChamRepository(t *testing.T) {
svc, repo, _ := newService(t)
in := validRegistration()
in.Password = "123" // quá ngắn
_, err := svc.Register(context.Background(), in)
assert.ErrorIs(t, err, domain.ErrInvalidInput)
assert.Empty(t, repo.byEmail)
}
func TestDangNhapSaiMatKhauVaEmailKhongTonTaiTraCungMotLoi(t *testing.T) {
svc, _, _ := newService(t)
_, err := svc.Register(context.Background(), validRegistration())
require.NoError(t, err)
_, errSaiMatKhau := svc.Login(context.Background(), "hhoang@example.com", "sai-mat-khau")
_, errKhongCoAi := svc.Login(context.Background(), "khong-ton-tai@example.com", "bat-ky")
// Hai nguyên nhân khác nhau, MỘT thông báo — để bên ngoài không
// dò được email nào đã đăng ký (§7.2).
assert.ErrorIs(t, errSaiMatKhau, domain.ErrInvalidCredentials)
assert.ErrorIs(t, errKhongCoAi, domain.ErrInvalidCredentials)
}
func TestDangNhapDungThiCapToken(t *testing.T) {
svc, _, _ := newService(t)
reg, err := svc.Register(context.Background(), validRegistration())
require.NoError(t, err)
res, err := svc.Login(context.Background(), "hhoang@example.com", "matkhau-du-dai-123")
require.NoError(t, err)
assert.Equal(t, reg.User.ID, res.User.ID)
assert.NotEmpty(t, res.Token)
}
Toàn bộ file này chạy trong khoảng một giây, không cần Docker, không cần PostgreSQL. Đó là phần thưởng cho việc khai báo phụ thuộc bằng interface ở §13.1.
17.4 Test chặn rò rỉ dữ liệu nhạy cảm
Thêm vào cùng file:
// Test này canh giữ quy tắc ở §10.1 thay cho bạn.
// Ngày có người thêm PasswordHash vào v1.UserRegistered, nó đỏ ngay.
func TestPayloadEventKhongChuaDuLieuNhayCam(t *testing.T) {
svc, _, writer := newService(t)
_, err := svc.Register(context.Background(), validRegistration())
require.NoError(t, err)
require.Len(t, writer.written, 1)
raw := string(writer.written[0].Payload)
assert.NotContains(t, raw, "matkhau-du-dai-123", "mật khẩu gốc lọt vào event")
assert.NotContains(t, raw, "$2a$", "hash bcrypt lọt vào event")
assert.NotContains(t, raw, "password", "trường liên quan mật khẩu lọt vào event")
}
Ba dòng assert, và chúng bảo vệ một thứ mà code review rất dễ bỏ sót — vì thêm một trường vào struct là thay đổi trông vô hại nhất trên đời.
17.5 Test tích hợp — user và outbox trong cùng transaction
Tạo internal/modules/identity/identity_integration_test.go:
package identity_test
import (
"context"
"os"
"testing"
"time"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
v1 "github.com/yourname/community/internal/contracts/v1"
"github.com/yourname/community/internal/modules/identity/domain"
"github.com/yourname/community/internal/modules/identity/repository"
"github.com/yourname/community/internal/modules/identity/service"
"github.com/yourname/community/internal/platform/logger"
"github.com/yourname/community/internal/platform/outbox"
"github.com/yourname/community/internal/platform/postgres"
"github.com/yourname/community/internal/platform/token"
)
func setup(t *testing.T) (context.Context, *service.Service, *pgxpool.Pool) {
t.Helper()
dsn := os.Getenv("TEST_DATABASE_URL")
if dsn == "" {
dsn = "postgres://app:secret@localhost:5432/community?sslmode=disable"
}
ctx := context.Background()
pool, err := postgres.Connect(ctx, postgres.DefaultConfig(dsn, "community-test"))
if err != nil {
t.Skipf("bỏ qua: cần PostgreSQL đang chạy (%v)", err)
}
t.Cleanup(pool.Close)
tokens, err := token.NewIssuer("secret-du-32-ky-tu-danh-cho-test!!!", time.Hour, "community")
require.NoError(t, err)
svc := service.New(
postgres.NewTxManager(pool),
repository.NewPostgres(pool),
outbox.NewWriter(),
tokens,
logger.New("error"),
)
return ctx, svc, pool
}
func randomRegistration() domain.Registration {
suffix := uuid.NewString()[:8]
return domain.Registration{
Name: "Người Dùng Thử",
Username: "test_" + suffix,
Email: "test_" + suffix + "@example.com",
Password: "matkhau-du-dai-123",
}
}
func TestDangKyGhiCaUserVaOutbox(t *testing.T) {
ctx, svc, pool := setup(t)
in := randomRegistration()
res, err := svc.Register(ctx, in)
require.NoError(t, err)
t.Cleanup(func() {
bg := context.Background()
_, _ = pool.Exec(bg, `DELETE FROM outbox WHERE aggregate_id = $1`, res.User.ID.String())
_, _ = pool.Exec(bg, `DELETE FROM users WHERE id = $1`, res.User.ID)
})
var users int
require.NoError(t, pool.QueryRow(ctx,
`SELECT count(*) FROM users WHERE id = $1`, res.User.ID).Scan(&users))
assert.Equal(t, 1, users)
var eventType string
var publishedAt *time.Time
require.NoError(t, pool.QueryRow(ctx, `
SELECT event_type, published_at FROM outbox
WHERE aggregate_id = $1`, res.User.ID.String()).Scan(&eventType, &publishedAt))
assert.Equal(t, v1.TypeUserRegistered, eventType)
assert.Nil(t, publishedAt, "relay chưa chạy nên event phải còn chờ")
}
// Test có giá trị nhất của bước này.
func TestDangKyTrungEmailKhongDeLaiGiTrongDatabase(t *testing.T) {
ctx, svc, pool := setup(t)
first := randomRegistration()
created, err := svc.Register(ctx, first)
require.NoError(t, err)
t.Cleanup(func() {
bg := context.Background()
_, _ = pool.Exec(bg, `DELETE FROM outbox WHERE aggregate_id = $1`, created.User.ID.String())
_, _ = pool.Exec(bg, `DELETE FROM users WHERE id = $1`, created.User.ID)
})
// Cùng email, username khác
second := randomRegistration()
second.Email = first.Email
_, err = svc.Register(ctx, second)
assert.ErrorIs(t, err, domain.ErrEmailTaken)
// Lần thất bại KHÔNG được để lại user, và cũng không để lại event.
var users, events int
require.NoError(t, pool.QueryRow(ctx,
`SELECT count(*) FROM users WHERE username = $1`, second.Username).Scan(&users))
require.NoError(t, pool.QueryRow(ctx,
`SELECT count(*) FROM outbox WHERE payload->>'username' = $1`, second.Username).Scan(&events))
assert.Equal(t, 0, users)
assert.Equal(t, 0, events, "transaction rollback phải kéo theo cả outbox")
}
func TestDangKyTrungUsernameBaoDungLoi(t *testing.T) {
ctx, svc, pool := setup(t)
first := randomRegistration()
created, err := svc.Register(ctx, first)
require.NoError(t, err)
t.Cleanup(func() {
bg := context.Background()
_, _ = pool.Exec(bg, `DELETE FROM outbox WHERE aggregate_id = $1`, created.User.ID.String())
_, _ = pool.Exec(bg, `DELETE FROM users WHERE id = $1`, created.User.ID)
})
second := randomRegistration()
second.Username = first.Username
_, err = svc.Register(ctx, second)
// Nếu chỗ này nhận ErrEmailTaken, hằng số tên ràng buộc ở §12.2 sai.
assert.ErrorIs(t, err, domain.ErrUsernameTaken)
}
Chạy tất cả:
go test ./internal/... -race -count=1
ok github.com/yourname/community/internal/contracts 0.28s
ok github.com/yourname/community/internal/modules/identity 1.84s
ok github.com/yourname/community/internal/modules/identity/service 1.61s
ok github.com/yourname/community/internal/platform/eventbus/inmem 0.26s
ok github.com/yourname/community/internal/platform/outbox 1.09s
ok github.com/yourname/community/internal/platform/password 3.42s
ok github.com/yourname/community/internal/platform/postgres 0.94s
ok github.com/yourname/community/internal/platform/token 0.31s
passwordmất 3.4 giây là bình thường — mỗi lầnHashtốn 250ms và test gọi nó hơn chục lần. Đó chính là thứ bảo vệ mật khẩu người dùng, nên đừng hạCostđể test chạy nhanh hơn.
18. Ai tạo dòng user_stats?
Người dùng vừa đăng ký cần một dòng trong user_stats với đủ 8 counter bằng 0. Câu hỏi: ai ghi dòng đó?
18.1 Cách sai, và vì sao nó hấp dẫn
// ❌ Bên trong service.Register
err = s.tx.Do(ctx, func(tx pgx.Tx) error {
saved, _ := s.repo.Insert(ctx, tx, user)
statsRepo.EnsureUserStats(ctx, tx, saved.ID) // ← identity ghi vào bảng của stats
...
})
Nó hấp dẫn vì đúng về mặt dữ liệu: cùng một transaction, không có khoảnh khắc nào user tồn tại mà thống kê thì chưa.
Nhưng nó vi phạm thẳng Luật 1. Hệ quả cụ thể, không phải lý thuyết:
identityphải importmodules/stats/repository→ hai module dính chặt vào nhau- Ngày
statsthêm counter thứ chín,identityphải sửa theo — dù nó chẳng liên quan gì - Ngày bạn muốn tách
statsthành service riêng (§15.2 của ARCHITECTURE),identitykhông còn ghi thẳng vào bảng đó được nữa, và đoạn code này phải viết lại từ đầu
Và bảng user_stats giờ có hai module cùng ghi. Đến khi một con số sai, câu hỏi "ai đã ghi dòng này?" không còn câu trả lời duy nhất.
18.2 Cách đúng
stats tự tạo dòng của mình khi nghe identity.user.registered.v1:
identity.Register
└─ INSERT users + INSERT outbox (một transaction)
│
└─ relay → Kafka topic "identity"
│
└─ stats.OnUserRegistered
└─ INSERT user_stats (user_id) ON CONFLICT DO NOTHING
Query EnsureUserStats viết ở bước 3 §9.5 chính là hàm này — nó đã sẵn sàng, chỉ chờ consumer ở bước 6.
Ba câu hỏi thường gặp về cách này:
"Khoá ngoại thì sao — user phải tồn tại trước?" Đúng, và nó luôn thoả mãn. Event chỉ rời khỏi bảng outbox sau khi transaction đã commit, nên tới lúc consumer chạy thì dòng users chắc chắn đã có.
"Có khoảng thời gian user tồn tại mà chưa có thống kê không?" Có — thường là dưới một giây. Với dữ liệu thống kê thì hoàn toàn chấp nhận được; đây đúng là loại dữ liệu mà nhất quán cuối cùng sinh ra để phục vụ. Nếu profile được mở đúng trong khoảnh khắc đó, tầng đọc hiển thị số 0 — đúng bằng giá trị mà dòng kia sắp có.
"Nếu consumer chết trước khi ghi thì sao?" Event chưa được commit offset, Kafka giao lại. Còn nếu nó được giao hai lần, ON CONFLICT DO NOTHING khiến lần thứ hai thành vô hại. Đó chính là lý do query được viết dưới dạng INSERT ... ON CONFLICT ngay từ đầu.
18.3 Đính chính tài liệu bước 3
Bước 3 §4.3 ban đầu chọn phương án ngược lại — tạo dòng user_stats trong cùng transaction với users, ở module identity. Khi viết Register thật thì mâu thuẫn lộ ra: không có cách nào làm điều đó mà không để identity import stats.
Tôi đã sửa lại §4.3 của bước 3 cho khớp với quyết định ở đây. Ghi lại chỗ này thay vì sửa lặng lẽ, vì bản thân sự việc là một bài học đáng giá: ranh giới module chỉ thật sự được kiểm chứng khi bạn viết dòng code đầu tiên đi ngang qua nó. Sơ đồ trên giấy không phát hiện được điều này.
19. Những thứ cố ý chưa có
| Chưa có | Vì sao | Sẽ làm ở |
|---|---|---|
outbox/relay.go |
Event đang nằm trong bảng, chưa ai đẩy lên Kafka | Bước 5 |
eventbus/kafka/ |
Cần relay trước mới có thứ để publish | Bước 5 |
idempotency/guard.go |
Module này chưa nghe event nào. Viết guard bây giờ là code chết không test nào chạm tới | Bước 6 (cùng consumer đầu tiên) |
stats.OnUserRegistered |
Cần Kafka chạy trước (§18.2) | Bước 6 |
| Refresh token, đăng xuất | Xem §8.4 — quyết định có ý thức, không phải bỏ sót | Khi cần |
PUT /me (đổi profile) |
Query UpdateUserProfile đã có từ bước 3, chỉ thiếu handler |
Bài tập tốt để tự làm |
| Giới hạn số lần đăng nhập sai | Cần Redis hoặc bảng đếm; đáng làm trước khi mở ra Internet | Trước khi lên production |
| Xác thực email | Cần dịch vụ gửi mail | Sau |
Về
idempotency/guard.go: bước 3 §13 có ghi nó thuộc bước 4. Tôi dời sang bước 6 vì bước này hoá ra không có consumer nào — và một lớp bọc không ai gọi thì không có gì chứng minh nó đúng.
20. Checklist hoàn thành
- [ ]
go build ./...thành công - [ ]
go vet ./...không cảnh báo - [ ]
go test ./internal/... -race -count=1— tất cảok, khôngSKIP - [ ]
go run ./cmd/apikhởi động, log ghiapi đang lắng nghe - [ ] Bỏ
JWT_SECRETkhỏi môi trường → khởi động thất bại kèm thông báo rõ ràng - [ ]
POST /auth/registertrả201kèm token - [ ] Response không chứa
passwordhaypassword_hash - [ ] Bảng
outboxcó đúng một dòngidentity.user.registered.v1,published_atlàNULL - [ ]
correlation_idtrong outbox khớp headerX-Request-IDcủa response - [ ] Đăng ký lại cùng email →
409, và outbox vẫn chỉ có một dòng - [ ] Đăng nhập bằng email viết HOA vẫn thành công
- [ ] Sai mật khẩu và email không tồn tại trả về cùng một thông báo, mất cùng một khoảng thời gian
- [ ]
GET /mekhông kèm token →401 - [ ]
GET /mevới token hợp lệ → thông tin đúng người dùng - [ ] Ctrl+C → log
nhận tín hiệu dừng, tiến trình thoát sạch
Cấu trúc thư mục lúc này:
project/
├── cmd/
│ ├── api/main.go ← MỚI
│ └── smoketest/main.go
├── db/migrations/ (không đổi — xem §1.2)
├── deployments/docker-compose.yml
├── docs/
│ ├── ARCHITECTURE.md
│ ├── 01-khoi-tao-du-an-va-ha-tang.md
│ ├── 02-nen-mong-platform.md
│ ├── 03-migrations-va-sqlc.md
│ └── 04-module-identity.md ← file này
├── internal/
│ ├── contracts/
│ │ ├── envelope.go
│ │ ├── topics.go
│ │ └── v1/{doc.go, post_events.go, user_events.go ← MỚI}
│ ├── modules/
│ │ ├── identity/ ← MỚI (trừ repository/)
│ │ │ ├── module.go
│ │ │ ├── domain/{user.go, errors.go}
│ │ │ ├── service/{service.go, register.go, login.go, service_test.go}
│ │ │ ├── repository/{query.sql, postgres.go, gen/}
│ │ │ ├── transport/http/{handler.go, router.go, dto.go}
│ │ │ └── identity_integration_test.go
│ │ └── stats/repository/
│ └── platform/
│ ├── config/ ← MỚI
│ ├── logger/ ← MỚI
│ ├── requestid/ ← MỚI
│ ├── httpx/ ← MỚI
│ ├── password/ ← MỚI
│ ├── token/ ← MỚI
│ ├── outbox/{query.sql, writer.go ← MỚI, gen/}
│ ├── idempotency/
│ ├── eventbus/
│ └── postgres/
├── sqlc.yaml
├── go.mod
└── go.sum
Commit:
git add .
git commit -m "feat(identity): đăng ký, đăng nhập, JWT và event UserRegistered vào outbox"
Phụ lục — Khi nào thật sự cần migration mới
§1.2 nói bước này không cần migration. Đây là mẫu cho ngày bạn nâng lên refresh token (giải pháp cho vấn đề ở §8.4).
Chưa cần tạo file này bây giờ — nó ở đây để bạn thấy một migration "liên quan tới auth" trông như thế nào.
db/migrations/000006_create_refresh_tokens.up.sql:
CREATE TABLE refresh_tokens (
-- Lưu HASH của token, không lưu token gốc. Lộ database thì kẻ tấn công
-- vẫn không có token dùng được — cùng lý lẽ với cột password_hash.
token_hash BYTEA PRIMARY KEY,
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
expires_at TIMESTAMPTZ NOT NULL,
revoked_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
user_agent TEXT,
ip INET
);
-- Cho màn hình "các thiết bị đang đăng nhập" và cho việc thu hồi
-- toàn bộ phiên của một người.
CREATE INDEX idx_refresh_tokens_user
ON refresh_tokens (user_id) WHERE revoked_at IS NULL;
-- Cho job dọn token hết hạn.
CREATE INDEX idx_refresh_tokens_expires_at ON refresh_tokens (expires_at);
db/migrations/000006_create_refresh_tokens.down.sql:
DROP TABLE IF EXISTS refresh_tokens;
Số hiệu là 000006 vì 000005 đã dành cho posts ở bước dựng module post.
Với bảng này, luồng auth đổi thành: access token TTL 15 phút (không thu hồi được, nhưng chết nhanh) + refresh token TTL 30 ngày (thu hồi được vì có trạng thái trong DB). Đăng xuất trở thành UPDATE refresh_tokens SET revoked_at = now().
Bước tiếp theo
Bước 5 — Outbox Relay và Kafka: viết outbox.Relay đọc bảng theo lô bằng FOR UPDATE SKIP LOCKED, cài eventbus/kafka bằng franz-go, và dựng cmd/relay.
Đó là bước dòng event bạn vừa tạo ở §16.3 rời khỏi PostgreSQL và xuất hiện trong Kafka UI ở localhost:8080 — lần đầu tiên hệ thống này thật sự trở thành event-driven.
All rights reserved