0

#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

400 hay 422? 400 nghĩa là "tôi không đọc nổi request này". 422 nghĩ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.go không phải import thêm package. Nếu thấy thừa, gọi thẳng middleware.Auth(h.verifier) trong Mount cũng đúng.


15. Lắp ráp: module.gocmd/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ế verifiertokenscù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)tokens xuấ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 main chỉ có ba dòng còn run làm hết?os.Exit không chạy các hàm defer. Nếu main vừa gọi os.Exit vừa defer pool.Close(), kết nối DB không bao giờ được đóng. Tách ra: run trả lỗi và đóng gọn mọi thứ, main chỉ 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

password mất 3.4 giây là bình thường — mỗi lần Hash tố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:

  • identity phải import modules/stats/repository → hai module dính chặt vào nhau
  • Ngày stats thêm counter thứ chín, identity phải sửa theo — dù nó chẳng liên quan gì
  • Ngày bạn muốn tách stats thành service riêng (§15.2 của ARCHITECTURE), identity khô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ông SKIP
  • [ ] go run ./cmd/api khởi động, log ghi api đang lắng nghe
  • [ ] Bỏ JWT_SECRET khỏi môi trường → khởi động thất bại kèm thông báo rõ ràng
  • [ ] POST /auth/register trả 201 kèm token
  • [ ] Response không chứa password hay password_hash
  • [ ] Bảng outbox có đúng một dòng identity.user.registered.v1, published_atNULL
  • [ ] correlation_id trong outbox khớp header X-Request-ID củ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 /me không kèm token → 401
  • [ ] GET /me vớ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à 000006000005 đã 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

Viblo
Hãy đăng ký một tài khoản Viblo để nhận được nhiều bài viết thú vị hơn.
Đăng kí