+1

#3 Bước 2 — Dựng nền móng Platform

Tài liệu thực hành. Mỗi package viết xong đều biên dịch và kiểm chứng được ngay trước khi sang package tiếp theo.

Yêu cầu: đã hoàn thành Bước 1 — hạ tầng Docker đang chạy Thời gian: 40–60 phút Kiến trúc tổng thể: ARCHITECTURE.md

Mục tiêu

Bước này viết ba package không chứa một dòng nghiệp vụ nào:

Package Vai trò
internal/platform/postgres Kết nối DB + helper chạy transaction
internal/contracts Envelope — định dạng chung của mọi event
internal/platform/eventbus Interface Publisher/Subscriber + bản in-memory để test

Kết thúc bước này bạn có bộ khung mà mọi module nghiệp vụ ở các bước sau sẽ dựa lên. Chưa có bảng dữ liệu, chưa có HTTP, chưa có Kafka thật.

Đây là bước dễ bị làm ẩu nhất. Nó không tạo ra tính năng nào nhìn thấy được, nên rất dễ viết qua loa cho xong. Nhưng TxManager sai một chỗ thì mọi transaction trong hệ thống sai theo, và bạn sẽ đi tìm nguyên nhân ở tầng nghiệp vụ — nơi không có lỗi gì cả.


Mục lục

  1. Platform là gì và không là gì
  2. Cài thư viện
  3. postgres — kết nối
  4. postgres — transaction
  5. contracts — Envelope
  6. eventbus — interface
  7. eventbus/inmem — bản cho test
  8. Kiểm chứng bằng test
  9. Nối tất cả vào smoketest
  10. Những thứ cố ý chưa có
  11. Checklist hoàn thành

1. Platform là gì và không là gì

Phép thử một câu

Nếu bạn mở một dự án Go hoàn toàn khác — một cửa hàng online, một hệ thống đặt vé — và copy package này sang, nó có chạy được không?

  • → đúng chỗ, để trong platform/
  • Không → nó thuộc về modules/

TxManager chạy được ở mọi dự án dùng PostgreSQL. IncrementPostsCount thì không — nó chỉ có nghĩa với dự án này. Ranh giới rõ ràng như vậy.

Vì sao phải nghiêm khắc

platform/ là tầng mà mọi module đều phụ thuộc vào. Một khi có nghiệp vụ lọt vào đó, hai chuyện xảy ra:

  1. Các module bắt đầu phụ thuộc gián tiếp vào nhau qua platform — Luật 1 bị phá mà không ai nhận ra, vì trên giấy tờ không module nào import module nào.
  2. Sửa một thứ trong platform có nguy cơ làm hỏng module không liên quan, nên dần dần không ai dám sửa gì cả.

Rule depguardARCHITECTURE.md §13.4 chặn platform/ import modules/. Nhưng nó không chặn được việc bạn chép nội dung nghiệp vụ vào platform. Cái đó chỉ có kỷ luật của bạn ngăn được.

Chiều phụ thuộc ở bước này

        eventbus ──────► contracts
                            ▲
        postgres            │
        (độc lập)           │
                         (chỉ dữ liệu,
                          không import gì
                          trong dự án)

contracts của đồ thị: nó không import bất kỳ package nội bộ nào. Đây là tính chất cần giữ bằng mọi giá — vì mọi module sẽ import nó, nếu contracts phụ thuộc ngược lại vào cái gì đó thì vòng lặp xuất hiện ngay.


2. Cài thư viện

go get github.com/jackc/pgx/v5
go get github.com/google/uuid
go get github.com/stretchr/testify
Thư viện Dùng để Ghi chú
pgx/v5 Driver Postgres Dùng thẳng, không qua database/sql
google/uuid Sinh event_id Nhỏ, không kéo theo dependency nào
testify Assert trong test Chỉ dùng ở file _test.go

Vì sao pgx mà không phải database/sql?

database/sql là interface chung cho mọi loại DB, nên nó chỉ phơi ra được phần giao của tất cả — mất hết tính năng riêng của Postgres. Dùng pgx trực tiếp bạn được: kiểu dữ liệu riêng của Postgres (jsonb, uuid, array) map thẳng sang Go, COPY để nạp dữ liệu hàng loạt, LISTEN/NOTIFY, và hiệu năng tốt hơn vì bỏ được một lớp trung gian.

sqlc ở bước sau cũng sinh code cho pgx/v5 — chọn đúng từ đầu thì về sau không phải viết lại.


3. postgres — kết nối

Tạo internal/platform/postgres/pool.go:

// Package postgres cung cấp kết nối PostgreSQL và helper transaction.
// Package này không chứa logic nghiệp vụ.
package postgres

import (
	"context"
	"fmt"
	"time"

	"github.com/jackc/pgx/v5/pgxpool"
)

// Config điều khiển hành vi của connection pool.
type Config struct {
	DSN string

	// MaxConns là số kết nối tối đa. Xem ghi chú về ngân sách kết nối bên dưới.
	MaxConns int32

	// MinConns là số kết nối giữ sẵn kể cả khi rảnh, tránh trả giá
	// bắt tay TCP + xác thực ở request đầu tiên sau lúc nhàn rỗi.
	MinConns int32

	// MaxConnLifetime buộc kết nối phải được tạo lại định kỳ.
	// Cần thiết khi có load balancer hoặc failover ở giữa.
	MaxConnLifetime time.Duration

	// MaxConnIdleTime đóng kết nối rảnh quá lâu để trả tài nguyên về server.
	MaxConnIdleTime time.Duration

	// AppName hiện trong pg_stat_activity — biết ai đang giữ kết nối.
	AppName string
}

// DefaultConfig trả về cấu hình hợp lý cho môi trường dev.
func DefaultConfig(dsn, appName string) Config {
	return Config{
		DSN:             dsn,
		MaxConns:        10,
		MinConns:        2,
		MaxConnLifetime: time.Hour,
		MaxConnIdleTime: 30 * time.Minute,
		AppName:         appName,
	}
}

// Connect mở pool và xác nhận server thực sự trả lời.
func Connect(ctx context.Context, cfg Config) (*pgxpool.Pool, error) {
	poolCfg, err := pgxpool.ParseConfig(cfg.DSN)
	if err != nil {
		return nil, fmt.Errorf("postgres: DSN không hợp lệ: %w", err)
	}

	poolCfg.MaxConns = cfg.MaxConns
	poolCfg.MinConns = cfg.MinConns
	poolCfg.MaxConnLifetime = cfg.MaxConnLifetime
	poolCfg.MaxConnIdleTime = cfg.MaxConnIdleTime

	if cfg.AppName != "" {
		poolCfg.ConnConfig.RuntimeParams["application_name"] = cfg.AppName
	}

	pool, err := pgxpool.NewWithConfig(ctx, poolCfg)
	if err != nil {
		return nil, fmt.Errorf("postgres: tạo pool thất bại: %w", err)
	}

	// pgxpool.NewWithConfig KHÔNG kết nối ngay — nó lazy.
	// Ping để lỗi cấu hình lộ ra lúc khởi động, không phải lúc có request đầu tiên.
	if err := pool.Ping(ctx); err != nil {
		pool.Close()
		return nil, fmt.Errorf("postgres: không kết nối được: %w", err)
	}

	return pool, nil
}

3.1 Vì sao phải Ping ngay

pgxpool.NewWithConfig chỉ dựng cấu trúc pool, không mở kết nối nào. Nếu bạn gõ sai mật khẩu hay Postgres chưa chạy, hàm này vẫn trả về nil error.

Không có Ping, tiến trình khởi động "thành công", log ghi "server started on :8000", rồi request đầu tiên của người dùng nhận lỗi 500. Bạn đi tìm bug ở handler trong khi lỗi thật nằm ở dòng cấu hình.

Fail fast lúc khởi động luôn tốt hơn fail muộn lúc phục vụ.

3.2 Ngân sách kết nối

PostgreSQL mặc định max_connections = 100. Con số đó chia cho toàn bộ hệ thống, không phải cho riêng một tiến trình.

Kiến trúc này có ba binary, mỗi binary một pool:

Tiến trình Số instance dự kiến MaxConns Tổng
api 3 10 30
worker 2 10 20
relay 1 5 5
Cộng 55

Còn dư khoảng 45 cho migration, psql thủ công, công cụ giám sát. Vừa đủ an toàn.

Cái bẫy: MaxConns thường bị copy giữa các service mà không ai cộng lại. Đến lúc scale api từ 3 lên 10 pod, tổng vọt lên 130 → Postgres từ chối kết nối mới → toàn hệ thống sập, kể cả những phần đang chạy tốt. Sự cố này rất hay gặp và triệu chứng thì trông như DB quá tải, dù DB đang rảnh.

Nhớ lại: MaxConns là ngân sách chung, không phải cấu hình riêng của từng service.

3.3 application_name — chi tiết nhỏ, giá trị lớn

Khi có sự cố, câu hỏi đầu tiên là "tiến trình nào đang giữ kết nối?". Có application_name thì trả lời được ngay:

SELECT application_name, state, count(*)
FROM pg_stat_activity
WHERE datname = 'community'
GROUP BY 1, 2 ORDER BY 3 DESC;
 application_name | state  | count
------------------+--------+-------
 community-api    | idle   |     8
 community-worker | active |     3
 community-relay  | idle   |     2

Không có nó, mọi dòng đều ghi application_name rỗng và bạn không phân biệt được gì.


4. postgres — transaction

Đây là phần quan trọng nhất của cả bước 2. Luật 3 — ghi dữ liệu và ghi outbox phải cùng transaction — được thực thi bởi đúng file này.

Tạo internal/platform/postgres/tx.go:

package postgres

import (
	"context"
	"errors"
	"fmt"

	"github.com/jackc/pgx/v5"
	"github.com/jackc/pgx/v5/pgxpool"
)

// TxManager chạy một hàm bên trong transaction.
type TxManager struct {
	pool *pgxpool.Pool
}

func NewTxManager(pool *pgxpool.Pool) *TxManager {
	return &TxManager{pool: pool}
}

// Do chạy fn trong một transaction.
//
//	fn trả nil   → COMMIT
//	fn trả error → ROLLBACK, lỗi được trả nguyên vẹn ra ngoài
//	fn panic     → ROLLBACK, panic được ném tiếp
//
// KHÔNG gọi Do lồng trong Do — xem ghi chú §4.3.
func (m *TxManager) Do(ctx context.Context, fn func(tx pgx.Tx) error) error {
	tx, err := m.pool.Begin(ctx)
	if err != nil {
		return fmt.Errorf("postgres: begin: %w", err)
	}

	// Rollback dùng context tách rời: nếu ctx đã bị huỷ (request timeout),
	// rollback với ctx đó sẽ thất bại và pgx phải huỷ luôn kết nối. Xem §4.2.
	rollbackCtx := context.WithoutCancel(ctx)

	defer func() {
		if p := recover(); p != nil {
			_ = tx.Rollback(rollbackCtx)
			panic(p) // ném tiếp để không nuốt mất panic
		}
	}()

	if err := fn(tx); err != nil {
		if rbErr := tx.Rollback(rollbackCtx); rbErr != nil && !errors.Is(rbErr, pgx.ErrTxClosed) {
			// Giữ CẢ HAI lỗi: lỗi nghiệp vụ vẫn là nguyên nhân gốc,
			// lỗi rollback là dấu hiệu kết nối đang có vấn đề.
			return errors.Join(err, fmt.Errorf("postgres: rollback: %w", rbErr))
		}
		return err
	}

	if err := tx.Commit(ctx); err != nil {
		return fmt.Errorf("postgres: commit: %w", err)
	}
	return nil
}

// DoValue giống Do nhưng trả thêm một giá trị.
// Là hàm ở cấp package vì Go không cho method có type parameter.
func DoValue[T any](ctx context.Context, m *TxManager, fn func(tx pgx.Tx) (T, error)) (T, error) {
	var out T
	err := m.Do(ctx, func(tx pgx.Tx) error {
		v, err := fn(tx)
		if err != nil {
			return err
		}
		out = v
		return nil
	})
	if err != nil {
		var zero T
		return zero, err
	}
	return out, nil
}

4.1 Vì sao không dùng defer tx.Rollback(ctx)

Bạn sẽ thấy cách viết này ở rất nhiều nơi:

tx, _ := pool.Begin(ctx)
defer tx.Rollback(ctx)   // ❌ trông gọn nhưng có ba vấn đề
// ...
return tx.Commit(ctx)

Nó dựa vào việc rollback sau commit là no-op (trả pgx.ErrTxClosed). Ba vấn đề:

  1. Lỗi rollback bị nuốt hoàn toàn. Rollback thất bại là dấu hiệu kết nối hỏng — thứ bạn rất muốn biết, nhưng defer đã ném nó đi.
  2. Vẫn dùng ctx đã có thể bị huỷ. Xem §4.2 ngay dưới.
  3. Mọi transaction trong dự án phải tự nhớ viết đúng. Quên một chỗ là rò rỉ một transaction — mà lỗi này không báo gì, chỉ âm thầm giữ kết nối tới khi hết timeout.

TxManager.Do gom toàn bộ chuyện đó vào một chỗ, viết đúng một lần. Người dùng nó chỉ cần trả error như bình thường.

4.2 context.WithoutCancel — chi tiết dễ bỏ sót

Kịch bản: request HTTP có timeout 5 giây. Transaction chạy tới giây thứ 5 thì ctx bị huỷ.

tx.Rollback(ctx)   // ctx đã Done → pgx không gửi được lệnh ROLLBACK

Khi rollback không gửi được, pgx không thể biết kết nối còn ở trạng thái nào, nên nó huỷ luôn kết nối đó thay vì trả về pool. Bình thường thì không sao. Nhưng đúng lúc tải cao — cũng chính là lúc timeout xảy ra nhiều nhất — bạn mất kết nối liên tục và pool phải mở lại từ đầu. Hệ thống đang chậm lại càng chậm thêm, theo vòng xoáy.

context.WithoutCancel(ctx) (có từ Go 1.21) tạo context giữ nguyên các giá trị nhưng bỏ tín hiệu huỷ. Rollback luôn gửi được, kết nối luôn về pool sạch sẽ.

4.3 Không gọi lồng nhau

// ❌ DEADLOCK
txm.Do(ctx, func(tx pgx.Tx) error {
    return txm.Do(ctx, func(tx2 pgx.Tx) error {  // lấy kết nối THỨ HAI
        ...
    })
})

Do bên trong xin một kết nối khác từ pool. Hai transaction trên hai kết nối khác nhau, không biết gì về nhau — nếu chúng đụng cùng một dòng dữ liệu, cái trong chờ cái ngoài nhả khoá, mà cái ngoài lại đang chờ cái trong xong. Khoá chết.

Quy tắc: Do chỉ được gọi ở tầng service, đúng một lần cho mỗi thao tác nghiệp vụ. Repository và outbox writer nhận tx pgx.Tx làm tham số, không bao giờ tự mở transaction.

Đây chính là lý do chữ ký hàm ở ARCHITECTURE.md §7.3 bắt buộc truyền tx:

func (w *Writer) Publish(ctx context.Context, tx pgx.Tx, ...) error

Không có tx thì không gọi được hàm. Trình biên dịch giữ luật thay cho bạn — đây là cách ép quy tắc đáng tin hơn mọi dòng comment.

4.4 Bắt lỗi nghiệp vụ từ Postgres

Repository ở các bước sau sẽ cần phân biệt "trùng email" với "DB hỏng". pgx phơi ra mã lỗi chuẩn của Postgres:

// internal/platform/postgres/errors.go
package postgres

import (
	"errors"

	"github.com/jackc/pgx/v5/pgconn"
)

// Mã lỗi Postgres: https://www.postgresql.org/docs/current/errcodes-appendix.html
const (
	CodeUniqueViolation     = "23505"
	CodeForeignKeyViolation = "23503"
	CodeSerializationFailure = "40001"
	CodeDeadlockDetected     = "40P01"
)

func errCode(err error) string {
	var pgErr *pgconn.PgError
	if errors.As(err, &pgErr) {
		return pgErr.Code
	}
	return ""
}

// IsUniqueViolation cho biết lỗi có phải do vi phạm UNIQUE không.
func IsUniqueViolation(err error) bool { return errCode(err) == CodeUniqueViolation }

// IsForeignKeyViolation cho biết lỗi có phải do vi phạm khoá ngoại không.
func IsForeignKeyViolation(err error) bool { return errCode(err) == CodeForeignKeyViolation }

// IsRetryable cho biết giao dịch có thể thử lại nguyên vẹn hay không.
func IsRetryable(err error) bool {
	c := errCode(err)
	return c == CodeSerializationFailure || c == CodeDeadlockDetected
}

Nhờ đó repository chuyển lỗi kỹ thuật thành lỗi nghiệp vụ:

if postgres.IsUniqueViolation(err) {
    return domain.ErrEmailTaken   // service hiểu; HTTP trả 409 thay vì 500
}

Không so khớp bằng chuỗi text của lỗi. Text thay đổi theo phiên bản Postgres và theo cả ngôn ngữ hệ thống; mã lỗi thì cố định.


5. contracts — Envelope

5.1 Envelope

Tạo internal/contracts/envelope.go:

// Package contracts định nghĩa hợp đồng event dùng chung cho toàn hệ thống.
//
// Package này KHÔNG import bất kỳ package nội bộ nào. Nó là lá của đồ thị
// phụ thuộc, nhờ vậy mọi module import được nó mà không tạo ra vòng lặp.
package contracts

import (
	"encoding/json"
	"errors"
	"fmt"
	"time"

	"github.com/google/uuid"
)

// Envelope là lớp bao ngoài của mọi event trong hệ thống.
type Envelope struct {
	// EventID định danh duy nhất một lần phát. Consumer dùng nó để
	// phát hiện message trùng (xem ARCHITECTURE.md §8).
	EventID string `json:"event_id"`

	// EventType theo mẫu <domain>.<entity>.<hành động quá khứ>.<version>
	// Ví dụ: "post.post.created.v1"
	EventType string `json:"event_type"`

	// AggregateID là ID của thực thể mà event nói về. Được dùng làm
	// Kafka partition key, nên mọi event của cùng một thực thể rơi vào
	// cùng partition và được xử lý ĐÚNG THỨ TỰ.
	AggregateID string `json:"aggregate_id"`

	// OccurredAt là lúc SỰ VIỆC xảy ra, không phải lúc publish.
	// Hai mốc này lệch nhau khi relay bị chậm.
	OccurredAt time.Time `json:"occurred_at"`

	// CorrelationID xuyên suốt một request: HTTP → outbox → consumer.
	// Có nó thì grep một ID là thấy toàn bộ hành trình.
	CorrelationID string `json:"correlation_id,omitempty"`

	// Payload là dữ liệu riêng của từng loại event, giữ nguyên dạng
	// JSON thô để envelope không phụ thuộc vào bất kỳ kiểu cụ thể nào.
	Payload json.RawMessage `json:"payload"`
}

var ErrInvalidEnvelope = errors.New("envelope không hợp lệ")

// NewEnvelope tạo envelope với EventID mới và mốc thời gian hiện tại.
func NewEnvelope(eventType, aggregateID string, payload any) (Envelope, error) {
	body, err := json.Marshal(payload)
	if err != nil {
		return Envelope{}, fmt.Errorf("contracts: mã hoá payload %s: %w", eventType, err)
	}
	e := Envelope{
		EventID:     uuid.NewString(),
		EventType:   eventType,
		AggregateID: aggregateID,
		OccurredAt:  time.Now().UTC(),
		Payload:     body,
	}
	if err := e.Validate(); err != nil {
		return Envelope{}, err
	}
	return e, nil
}

// WithCorrelation trả về bản sao có gắn correlation ID.
func (e Envelope) WithCorrelation(id string) Envelope {
	e.CorrelationID = id
	return e
}

// Validate kiểm tra các trường bắt buộc.
// Gọi ở BIÊN GIỚI: trước khi ghi outbox và sau khi nhận từ broker.
func (e Envelope) Validate() error {
	switch {
	case e.EventID == "":
		return fmt.Errorf("%w: thiếu event_id", ErrInvalidEnvelope)
	case e.EventType == "":
		return fmt.Errorf("%w: thiếu event_type", ErrInvalidEnvelope)
	case e.AggregateID == "":
		return fmt.Errorf("%w: thiếu aggregate_id (%s)", ErrInvalidEnvelope, e.EventType)
	case e.OccurredAt.IsZero():
		return fmt.Errorf("%w: thiếu occurred_at (%s)", ErrInvalidEnvelope, e.EventType)
	case len(e.Payload) == 0:
		return fmt.Errorf("%w: payload rỗng (%s)", ErrInvalidEnvelope, e.EventType)
	}
	return nil
}

// DecodePayload giải mã payload về kiểu cụ thể.
//
//	p, err := contracts.DecodePayload[v1.PostCreated](e)
func DecodePayload[T any](e Envelope) (T, error) {
	var out T
	if err := json.Unmarshal(e.Payload, &out); err != nil {
		return out, fmt.Errorf("contracts: giải mã payload %s: %w", e.EventType, err)
	}
	return out, nil
}

Ghi chú: contracts được phép có method

ARCHITECTURE.md §4.2 ban đầu viết "chỉ chứa struct dữ liệu và hằng số, không method". Quy tắc đó quá chặt và tôi đã nới lại cho đúng với ý định ban đầu:

Được phép: thao tác trên chính dữ liệu của nó — khởi tạo, kiểm tra hợp lệ, mã hoá/giải mã. Không được phép: truy cập DB, gọi mạng, quyết định nghiệp vụ, import bất kỳ package nội bộ nào.

Mục đích thật của quy tắc là giữ contracts không phụ thuộc vào ai, chứ không phải cấm method. Validate() không làm contracts phụ thuộc thêm thứ gì.

5.2 Vì sao Payloadjson.RawMessage

Nếu khai Payload any, envelope sẽ phải biết về mọi kiểu event trong hệ thống để giải mã đúng — tức là contracts phụ thuộc vào tất cả, đúng thứ ta đang tránh.

json.RawMessage giữ payload ở dạng byte JSON thô. Envelope vận chuyển nó mà không cần hiểu nội dung; chỉ consumer — người biết mình chờ event gì — mới giải mã:

p, err := contracts.DecodePayload[v1.PostCreated](e)

Đây cũng là điều làm hệ thống chịu được thay đổi: thêm một loại event mới không cần sửa một dòng nào trong envelope.go.

5.3 Topics

Tạo internal/contracts/topics.go:

package contracts

// Topic Kafka. Mỗi topic gom event của một domain.
//
// Gom theo domain (không phải mỗi loại event một topic) vì:
//   - Event trong cùng domain thường cần giữ đúng thứ tự với nhau
//   - Consumer đăng ký một lần thay vì hàng chục lần
//   - Số partition tổng không phình theo số loại event
const (
	TopicIdentity = "identity"
	TopicPost     = "post"
	TopicSocial   = "social"
)

// DLQ trả về tên Dead Letter Queue của một topic.
func DLQ(topic string) string { return topic + ".DLQ" }

5.4 Event phiên bản 1

Tạo internal/contracts/v1/doc.go:

// Package v1 chứa payload event phiên bản 1.
//
// QUY TẮC VERSION
//
// Được phép sửa trực tiếp trong v1 (tương thích ngược):
//   - Thêm trường mới CÓ giá trị zero hợp lệ
//
// Bắt buộc tạo v2 (phá vỡ tương thích):
//   - Xoá trường, đổi tên trường
//   - Đổi kiểu dữ liệu
//   - Đổi ngữ nghĩa của trường sẵn có
//
// Quy trình lên v2:
//   1. Tạo package v2 với type mới "<...>.v2"
//   2. Publisher phát CẢ HAI trong giai đoạn chuyển tiếp
//   3. Chuyển từng consumer sang v2
//   4. Không còn ai nghe v1 → ngừng phát v1
//
// KHÔNG BAO GIỜ đổi nghĩa của event đã phát hành: trong Kafka có thể
// còn message v1 chưa xử lý, và consumer đọc chúng bằng schema mới sẽ
// hiểu sai dữ liệu — âm thầm, không báo lỗi.
package v1

Tạo internal/contracts/v1/post_events.go:

package v1

import "time"

const (
	TypePostCreated = "post.post.created.v1"
	TypePostViewed  = "post.post.viewed.v1"
	TypePostDeleted = "post.post.deleted.v1"
)

type PostCreated struct {
	PostID    string    `json:"post_id"`
	AuthorID  string    `json:"author_id"`
	Title     string    `json:"title"`
	Slug      string    `json:"slug"`
	ReadTime  int       `json:"read_time"`
	CreatedAt time.Time `json:"created_at"`
}

type PostViewed struct {
	PostID string `json:"post_id"`
	// AuthorID có mặt ở đây để consumer stats cập nhật total_post_views
	// của tác giả mà KHÔNG phải query ngược sang bảng posts.
	// Trùng lặp dữ liệu là cái giá rẻ để giữ module độc lập.
	AuthorID string `json:"author_id"`
	ViewerID string `json:"viewer_id"`
}

type PostDeleted struct {
	PostID   string `json:"post_id"`
	AuthorID string `json:"author_id"`
}

Ba event này là đủ cho bước 2. Các event còn lại (identity, social) sẽ thêm khi dựng module tương ứng — thêm sớm chỉ tạo ra code chết không ai kiểm chứng.


6. eventbus — interface

Tạo internal/platform/eventbus/bus.go:

// Package eventbus định nghĩa giao diện phát/nhận event.
//
// Module nghiệp vụ CHỈ nhìn thấy các interface trong file này.
// Kafka, NATS hay in-memory đều nằm sau chúng, nên đổi hạ tầng
// không phải sửa một dòng nghiệp vụ nào.
package eventbus

import (
	"context"

	"github.com/yourname/community/internal/contracts"
)

// Publisher phát event lên một topic.
//
// Trong luồng chính, KHÔNG gọi Publisher trực tiếp từ service —
// hãy ghi vào outbox (ARCHITECTURE.md §7). Publisher này dành cho
// tiến trình relay, tiến trình duy nhất được nói chuyện với broker.
type Publisher interface {
	Publish(ctx context.Context, topic string, events ...contracts.Envelope) error
}

// Handler xử lý một event.
//
// Trả nil       → đánh dấu đã xử lý, commit offset
// Trả lỗi thường→ retry theo backoff
// Trả ErrPermanent → vào thẳng DLQ, không retry
type Handler func(ctx context.Context, e contracts.Envelope) error

// Subscriber đăng ký handler và chạy vòng lặp tiêu thụ.
type Subscriber interface {
	// Subscribe đăng ký handler cho (topic, group).
	// Mỗi group nhận MỘT BẢN SAO của mọi event trên topic;
	// trong cùng group, mỗi event chỉ được một instance xử lý.
	//
	// Gọi Subscribe trước khi gọi Run.
	Subscribe(topic, group string, h Handler)

	// Run chạy tới khi ctx bị huỷ. Chặn luồng gọi.
	Run(ctx context.Context) error
}

Vì sao Publish nhận biến số event?

Relay đọc outbox theo lô 100 dòng (ARCHITECTURE.md §7.4). Nếu chữ ký chỉ nhận một event, relay phải gọi 100 lần → 100 lượt đi về mạng cho việc lẽ ra chỉ cần một.

Dạng biến số không làm phức tạp trường hợp thường gặp: Publish(ctx, topic, ev) vẫn đọc y như cũ.

6.1 Phân loại lỗi

Tạo internal/platform/eventbus/errors.go:

package eventbus

import (
	"errors"
	"fmt"
)

// ErrPermanent đánh dấu lỗi mà retry chắc chắn vô ích:
// payload hỏng, thiếu trường bắt buộc, vi phạm ràng buộc nghiệp vụ.
// Event mang lỗi này đi thẳng vào DLQ.
var ErrPermanent = errors.New("permanent")

// Permanent bọc err thành lỗi vĩnh viễn.
func Permanent(err error) error {
	if err == nil {
		return nil
	}
	return fmt.Errorf("%w: %w", ErrPermanent, err)
}

// IsPermanent cho biết có nên bỏ qua retry hay không.
func IsPermanent(err error) bool { return errors.Is(err, ErrPermanent) }

Vì sao chỉ có ErrPermanent mà không có ErrTransient?

Vì mặc định an toàn hơn là retry. Lỗi không được đánh dấu gì sẽ được thử lại — mất DB một nhịp thì lần sau thành công. Nếu làm ngược lại (phải đánh dấu mới retry), một lỗi tạm thời bị quên đánh dấu sẽ khiến event mất luôn.

Chọn mặc định sao cho quên khai báo dẫn tới hậu quả nhẹ hơn. Ở đây quên là bị retry thừa vài lần; còn quên ở chiều kia là mất dữ liệu.

6.2 Router theo loại event

Tạo internal/platform/eventbus/router.go:

package eventbus

import (
	"context"

	"github.com/yourname/community/internal/contracts"
)

// Route gom nhiều handler theo EventType thành một Handler duy nhất.
//
//	bus.Subscribe(contracts.TopicPost, "stats-service", eventbus.Route(map[string]eventbus.Handler{
//	    v1.TypePostCreated: sub.OnPostCreated,
//	    v1.TypePostViewed:  sub.OnPostViewed,
//	}))
func Route(routes map[string]Handler) Handler {
	return func(ctx context.Context, e contracts.Envelope) error {
		h, ok := routes[e.EventType]
		if !ok {
			// Loại event không quan tâm → BỎ QUA, không phải lỗi. Xem §6.3.
			return nil
		}
		return h(ctx, e)
	}
}

6.3 Vì sao event lạ phải bỏ qua, không được báo lỗi

Đây là quyết định nhỏ nhưng hậu quả lớn.

Topic post chở nhiều loại event. Consumer stats chỉ quan tâm ba loại. Ngày mai module post thêm post.post.pinned.v1.

Nếu Route trả lỗi cho loại lạ:

[stats-service] ERROR: unknown event type post.post.pinned.v1
[stats-service] retry 1... retry 2... retry 3...
[stats-service] → DLQ

Mọi consumer hiện có đều gãy chỉ vì một module khác thêm event mới. Đúng thứ Event-Driven sinh ra để tránh. Tệ hơn: sự cố xảy ra ở stats, còn nguyên nhân nằm ở post — không ai nối được hai đầu.

Bỏ qua trong im lặng giữ đúng lời hứa của kiến trúc: thêm event mới không bao giờ làm hỏng consumer cũ.

Bù lại, ở môi trường dev nên log mức debug cho loại event bị bỏ qua — để khi bạn viết handler mà mãi không thấy nó chạy, log nói ngay rằng event đã tới nhưng không khớp EventType nào (thường là do gõ sai hằng số).


7. eventbus/inmem — bản cho test

Tạo internal/platform/eventbus/inmem/bus.go:

// Package inmem cài eventbus hoàn toàn trong bộ nhớ, dành cho unit test.
//
// Khác biệt CÓ CHỦ ĐÍCH so với Kafka:
//   - Giao event ĐỒNG BỘ ngay trong lời gọi Publish → test tất định,
//     không cần sleep, không cần chờ.
//   - Không bền: mất hết khi tiến trình kết thúc.
//   - Không retry, không DLQ: lỗi từ handler được trả thẳng cho Publish.
package inmem

import (
	"context"
	"sync"

	"github.com/yourname/community/internal/contracts"
	"github.com/yourname/community/internal/platform/eventbus"
)

type subscription struct {
	topic   string
	group   string
	handler eventbus.Handler
}

// Published ghi lại một lần phát, phục vụ assert trong test.
type Published struct {
	Topic string
	Event contracts.Envelope
}

type Bus struct {
	mu     sync.Mutex
	subs   []subscription
	record []Published
}

func New() *Bus { return &Bus{} }

// Bảo đảm lúc biên dịch rằng Bus thoả cả hai interface.
var (
	_ eventbus.Publisher  = (*Bus)(nil)
	_ eventbus.Subscriber = (*Bus)(nil)
)

func (b *Bus) Subscribe(topic, group string, h eventbus.Handler) {
	b.mu.Lock()
	defer b.mu.Unlock()
	b.subs = append(b.subs, subscription{topic: topic, group: group, handler: h})
}

func (b *Bus) Publish(ctx context.Context, topic string, events ...contracts.Envelope) error {
	for _, e := range events {
		if err := e.Validate(); err != nil {
			return err // bắt envelope hỏng ngay trong test, không đợi tới production
		}

		b.mu.Lock()
		b.record = append(b.record, Published{Topic: topic, Event: e})
		targets := make([]eventbus.Handler, 0, len(b.subs))
		for _, s := range b.subs {
			if s.topic == topic {
				targets = append(targets, s.handler)
			}
		}
		b.mu.Unlock() // nhả khoá TRƯỚC khi gọi handler, tránh deadlock nếu handler publish tiếp

		for _, h := range targets {
			if err := h(ctx, e); err != nil {
				return err
			}
		}
	}
	return nil
}

// Run chặn tới khi ctx bị huỷ, để Bus dùng thay Subscriber thật được.
func (b *Bus) Run(ctx context.Context) error {
	<-ctx.Done()
	return nil
}

// Published trả về mọi event đã phát, theo đúng thứ tự.
func (b *Bus) Published() []Published {
	b.mu.Lock()
	defer b.mu.Unlock()
	return append([]Published(nil), b.record...)
}

// PublishedOn lọc theo topic.
func (b *Bus) PublishedOn(topic string) []contracts.Envelope {
	b.mu.Lock()
	defer b.mu.Unlock()
	var out []contracts.Envelope
	for _, p := range b.record {
		if p.Topic == topic {
			out = append(out, p.Event)
		}
	}
	return out
}

// Reset xoá lịch sử, giữ nguyên các đăng ký.
func (b *Bus) Reset() {
	b.mu.Lock()
	defer b.mu.Unlock()
	b.record = nil
}

7.1 Hai chi tiết đáng chú ý

Nhả khoá trước khi gọi handler. Nếu giữ b.mu trong lúc chạy handler, mà handler đó lại gọi Publish (chuyện rất bình thường: xử lý event A sinh ra event B), nó sẽ tự khoá chính mình. sync.Mutex của Go không cho phép khoá lại (không reentrant) → treo vĩnh viễn, không có thông báo gì.

Validate() ngay trong Publish. Bản in-memory nghiêm khắc hơn cả bản thật, và đó là điều bạn muốn: một envelope thiếu AggregateID sẽ làm test đỏ trên máy bạn, thay vì trở thành message rơi vào DLQ trên production lúc 2 giờ sáng.

7.2 Nó không giả lập điều gì

Bản in-memory không mô phỏng: giao trùng message, giao sai thứ tự, độ trễ mạng, consumer rebalance.

Nghĩa là test với inmem không chứng minh được consumer của bạn idempotent. Việc đó cần Postgres thật với bảng consumed_events — sẽ làm ở bước sau, khi bảng đó tồn tại.

Biết rõ một công cụ không kiểm chứng điều gì cũng quan trọng ngang việc biết nó kiểm chứng điều gì.


8. Kiểm chứng bằng test

8.1 Test contracts

Tạo internal/contracts/envelope_test.go:

package contracts_test

import (
	"testing"

	"github.com/stretchr/testify/assert"
	"github.com/stretchr/testify/require"

	"github.com/yourname/community/internal/contracts"
	v1 "github.com/yourname/community/internal/contracts/v1"
)

func TestNewEnvelope_DienDuTruongBatBuoc(t *testing.T) {
	e, err := contracts.NewEnvelope(v1.TypePostCreated, "post-123", v1.PostCreated{
		PostID:   "post-123",
		AuthorID: "user-1",
		Title:    "Xin chào",
	})
	require.NoError(t, err)

	assert.NotEmpty(t, e.EventID)
	assert.Equal(t, v1.TypePostCreated, e.EventType)
	assert.Equal(t, "post-123", e.AggregateID)
	assert.False(t, e.OccurredAt.IsZero())
	require.NoError(t, e.Validate())
}

func TestNewEnvelope_MoiLanMotEventIDKhacNhau(t *testing.T) {
	a, err := contracts.NewEnvelope(v1.TypePostCreated, "p1", v1.PostCreated{})
	require.NoError(t, err)
	b, err := contracts.NewEnvelope(v1.TypePostCreated, "p1", v1.PostCreated{})
	require.NoError(t, err)

	// EventID phải khác nhau, nếu không cơ chế chống trùng ở consumer
	// sẽ nuốt mất event thứ hai.
	assert.NotEqual(t, a.EventID, b.EventID)
}

func TestValidate_BatThieuAggregateID(t *testing.T) {
	e, err := contracts.NewEnvelope(v1.TypePostCreated, "", v1.PostCreated{})
	require.Error(t, err)
	assert.ErrorIs(t, err, contracts.ErrInvalidEnvelope)
	assert.Empty(t, e.EventID)
}

func TestDecodePayload_KhopVoiLucTao(t *testing.T) {
	in := v1.PostCreated{PostID: "p1", AuthorID: "u1", Title: "Tiêu đề", ReadTime: 5}
	e, err := contracts.NewEnvelope(v1.TypePostCreated, in.PostID, in)
	require.NoError(t, err)

	out, err := contracts.DecodePayload[v1.PostCreated](e)
	require.NoError(t, err)
	assert.Equal(t, in.PostID, out.PostID)
	assert.Equal(t, in.Title, out.Title)
	assert.Equal(t, in.ReadTime, out.ReadTime)
}

Vì sao package test là contracts_test chứ không phải contracts?

Hậu tố _test buộc test chỉ dùng được API công khai — đúng góc nhìn của module sẽ import nó. Nếu test chạm được vào nội bộ, bạn có thể vô tình viết test xanh cho một API mà bên ngoài không dùng nổi.

8.2 Test eventbus

Tạo internal/platform/eventbus/inmem/bus_test.go:

package inmem_test

import (
	"context"
	"errors"
	"testing"

	"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/platform/eventbus"
	"github.com/yourname/community/internal/platform/eventbus/inmem"
)

func envelope(t *testing.T, eventType string) contracts.Envelope {
	t.Helper()
	e, err := contracts.NewEnvelope(eventType, "post-1", v1.PostCreated{PostID: "post-1"})
	require.NoError(t, err)
	return e
}

func TestBus_MoiGroupNhanMotBanSao(t *testing.T) {
	bus := inmem.New()
	var a, b int

	bus.Subscribe(contracts.TopicPost, "stats", func(context.Context, contracts.Envelope) error {
		a++
		return nil
	})
	bus.Subscribe(contracts.TopicPost, "search", func(context.Context, contracts.Envelope) error {
		b++
		return nil
	})

	require.NoError(t, bus.Publish(context.Background(), contracts.TopicPost, envelope(t, v1.TypePostCreated)))

	assert.Equal(t, 1, a)
	assert.Equal(t, 1, b)
}

func TestBus_KhongGiaoNhamTopic(t *testing.T) {
	bus := inmem.New()
	called := false
	bus.Subscribe(contracts.TopicSocial, "g", func(context.Context, contracts.Envelope) error {
		called = true
		return nil
	})

	require.NoError(t, bus.Publish(context.Background(), contracts.TopicPost, envelope(t, v1.TypePostCreated)))
	assert.False(t, called)
}

func TestRoute_BoQuaEventLa(t *testing.T) {
	handled := false
	h := eventbus.Route(map[string]eventbus.Handler{
		v1.TypePostCreated: func(context.Context, contracts.Envelope) error {
			handled = true
			return nil
		},
	})

	// Loại event không có trong bảng định tuyến: KHÔNG được coi là lỗi,
	// nếu không thì thêm event mới sẽ làm gãy mọi consumer đang chạy.
	err := h(context.Background(), envelope(t, "post.post.pinned.v1"))
	require.NoError(t, err)
	assert.False(t, handled)
}

func TestBus_TuChoiEnvelopeHong(t *testing.T) {
	bus := inmem.New()
	err := bus.Publish(context.Background(), contracts.TopicPost, contracts.Envelope{EventType: "x"})
	assert.ErrorIs(t, err, contracts.ErrInvalidEnvelope)
}

func TestPermanent_PhanBietDuocVoiLoiThuong(t *testing.T) {
	perm := eventbus.Permanent(errors.New("payload hỏng"))
	assert.True(t, eventbus.IsPermanent(perm))
	assert.False(t, eventbus.IsPermanent(errors.New("mất kết nối DB")))

	// Lỗi gốc vẫn phải đọc được sau khi bọc
	assert.Contains(t, perm.Error(), "payload hỏng")
}

8.3 Test TxManager — cần Postgres thật

Transaction không giả lập được cho ra hồn: rollback, isolation, khoá đều là hành vi của server. Ta dùng luôn Postgres đã dựng ở bước 1.

Tạo internal/platform/postgres/tx_test.go:

package postgres_test

import (
	"context"
	"errors"
	"os"
	"testing"

	"github.com/jackc/pgx/v5"
	"github.com/stretchr/testify/assert"
	"github.com/stretchr/testify/require"

	"github.com/yourname/community/internal/platform/postgres"
)

func setup(t *testing.T) (context.Context, *postgres.TxManager, func(string) int) {
	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: không kết nối được Postgres (%v) — chạy 'docker compose up -d' trước", err)
	}
	t.Cleanup(pool.Close)

	_, err = pool.Exec(ctx, `CREATE TABLE IF NOT EXISTS tx_probe (id TEXT PRIMARY KEY)`)
	require.NoError(t, err)
	_, err = pool.Exec(ctx, `TRUNCATE tx_probe`)
	require.NoError(t, err)
	t.Cleanup(func() { _, _ = pool.Exec(context.Background(), `DROP TABLE IF EXISTS tx_probe`) })

	count := func(id string) int {
		var n int
		require.NoError(t, pool.QueryRow(ctx, `SELECT count(*) FROM tx_probe WHERE id = $1`, id).Scan(&n))
		return n
	}
	return ctx, postgres.NewTxManager(pool), count
}

func TestDo_CommitKhiThanhCong(t *testing.T) {
	ctx, txm, count := setup(t)

	err := txm.Do(ctx, func(tx pgx.Tx) error {
		_, err := tx.Exec(ctx, `INSERT INTO tx_probe (id) VALUES ($1)`, "ok")
		return err
	})

	require.NoError(t, err)
	assert.Equal(t, 1, count("ok"))
}

func TestDo_RollbackKhiLoi(t *testing.T) {
	ctx, txm, count := setup(t)
	sentinel := errors.New("nghiệp vụ từ chối")

	err := txm.Do(ctx, func(tx pgx.Tx) error {
		if _, err := tx.Exec(ctx, `INSERT INTO tx_probe (id) VALUES ($1)`, "rollback"); err != nil {
			return err
		}
		return sentinel // ← đã ghi rồi mới lỗi
	})

	// Lỗi gốc phải tới tay người gọi nguyên vẹn, không bị bọc mất
	assert.ErrorIs(t, err, sentinel)
	// Và dòng vừa ghi phải biến mất
	assert.Equal(t, 0, count("rollback"))
}

func TestDo_RollbackKhiPanic(t *testing.T) {
	ctx, txm, count := setup(t)

	assert.Panics(t, func() {
		_ = txm.Do(ctx, func(tx pgx.Tx) error {
			_, err := tx.Exec(ctx, `INSERT INTO tx_probe (id) VALUES ($1)`, "panic")
			require.NoError(t, err)
			panic("bug ở tầng service")
		})
	})

	// Panic phải được ném tiếp (assert.Panics ở trên đã kiểm tra)
	// NHƯNG transaction vẫn phải rollback sạch sẽ.
	assert.Equal(t, 0, count("panic"))
}

func TestDo_RollbackDuocCaKhiContextBiHuy(t *testing.T) {
	ctx, txm, count := setup(t)
	cancelCtx, cancel := context.WithCancel(ctx)

	err := txm.Do(cancelCtx, func(tx pgx.Tx) error {
		if _, err := tx.Exec(cancelCtx, `INSERT INTO tx_probe (id) VALUES ($1)`, "cancelled"); err != nil {
			return err
		}
		cancel() // mô phỏng request timeout giữa chừng
		return errors.New("hết giờ")
	})

	require.Error(t, err)
	assert.Equal(t, 0, count("cancelled"))
}

func TestDoValue_TraGiaTriKhiCommit(t *testing.T) {
	ctx, txm, _ := setup(t)

	id, err := postgres.DoValue(ctx, txm, func(tx pgx.Tx) (string, error) {
		var out string
		err := tx.QueryRow(ctx, `INSERT INTO tx_probe (id) VALUES ($1) RETURNING id`, "value").Scan(&out)
		return out, err
	})

	require.NoError(t, err)
	assert.Equal(t, "value", id)
}

Test cuối là bằng chứng cho §4.2: nếu bỏ context.WithoutCancel, rollback sẽ thất bại vì context đã huỷ và test này đỏ.

8.4 Chạy

go build ./...
go test ./internal/... -v -race -count=1

Kết quả mong đợi:

ok  github.com/yourname/community/internal/contracts               0.31s
ok  github.com/yourname/community/internal/platform/eventbus/inmem 0.28s
ok  github.com/yourname/community/internal/platform/postgres       0.94s

Vì sao có -race? Hệ thống này chạy nhiều goroutine (consumer, relay, HTTP handler). Race condition là loại bug xuất hiện ngẫu nhiên và biến mất khi bạn thêm log vào để tìm nó. -race phát hiện được ngay từ lần chạy đầu. Bật nó trong CI ngay từ ngày đầu tiên.

Vì sao có -count=1? Go cache kết quả test. -count=1 tắt cache, buộc chạy thật — cần thiết khi test phụ thuộc trạng thái bên ngoài như DB.

Nếu test Postgres báo SKIP, hạ tầng chưa chạy:

docker compose -f deployments/docker-compose.yml up -d

9. Nối tất cả vào smoketest

Cập nhật cmd/smoketest/main.go để dùng chính các package vừa viết — vừa kiểm chứng chúng lắp vào nhau được, vừa cho thấy hình dáng code thật ở các bước sau:

package main

import (
	"context"
	"errors"
	"log"
	"os"
	"time"

	"github.com/jackc/pgx/v5"

	"github.com/yourname/community/internal/contracts"
	v1 "github.com/yourname/community/internal/contracts/v1"
	"github.com/yourname/community/internal/platform/eventbus"
	"github.com/yourname/community/internal/platform/eventbus/inmem"
	"github.com/yourname/community/internal/platform/postgres"
)

func env(key, def string) string {
	if v := os.Getenv(key); v != "" {
		return v
	}
	return def
}

func main() {
	ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
	defer cancel()

	// ── 1. Kết nối DB ─────────────────────────────────────────
	dsn := env("DATABASE_URL", "postgres://app:secret@localhost:5432/community?sslmode=disable")
	pool, err := postgres.Connect(ctx, postgres.DefaultConfig(dsn, "community-smoketest"))
	if err != nil {
		log.Fatalf("❌ %v", err)
	}
	defer pool.Close()
	log.Println("✅ PostgreSQL: kết nối OK")

	// ── 2. Transaction rollback đúng ──────────────────────────
	txm := postgres.NewTxManager(pool)
	sentinel := "cố ý lỗi"
	err = txm.Do(ctx, func(tx pgx.Tx) error {
		if _, err := tx.Exec(ctx, `SELECT 1`); err != nil {
			return err
		}
		return errors.New(sentinel)
	})
	if err == nil || err.Error() != sentinel {
		log.Fatalf("❌ TxManager: mong đợi lỗi %q, nhận %v", sentinel, err)
	}
	log.Println("✅ TxManager: rollback OK")

	// ── 3. Event bus giao đúng ────────────────────────────────
	bus := inmem.New()
	received := make(chan contracts.Envelope, 1)

	bus.Subscribe(contracts.TopicPost, "smoketest", eventbus.Route(map[string]eventbus.Handler{
		v1.TypePostCreated: func(_ context.Context, e contracts.Envelope) error {
			received <- e
			return nil
		},
	}))

	ev, err := contracts.NewEnvelope(v1.TypePostCreated, "post-1", v1.PostCreated{
		PostID:   "post-1",
		AuthorID: "user-1",
		Title:    "Bài viết đầu tiên",
	})
	if err != nil {
		log.Fatalf("❌ %v", err)
	}
	if err := bus.Publish(ctx, contracts.TopicPost, ev.WithCorrelation("smoke-001")); err != nil {
		log.Fatalf("❌ %v", err)
	}

	select {
	case got := <-received:
		p, err := contracts.DecodePayload[v1.PostCreated](got)
		if err != nil {
			log.Fatalf("❌ %v", err)
		}
		log.Printf("✅ EventBus: nhận %s → %q (correlation=%s)", got.EventType, p.Title, got.CorrelationID)
	default:
		log.Fatal("❌ EventBus: không nhận được event")
	}

	log.Println("🎉 Nền móng platform sẵn sàng.")
}

Chạy:

go run ./cmd/smoketest
✅ PostgreSQL: kết nối OK
✅ TxManager: rollback OK
✅ EventBus: nhận post.post.created.v1 → "Bài viết đầu tiên" (correlation=smoke-001)
🎉 Nền móng platform sẵn sàng.

10. Những thứ cố ý chưa có

Danh sách này quan trọng ngang phần đã làm — nó cho biết bạn không quên gì:

Chưa có Vì sao Sẽ làm ở
outbox/ Cần bảng outbox tồn tại trước Bước 3
idempotency/ Cần bảng consumed_events Bước 3
eventbus/kafka/ Chỉ có nghĩa khi đã có event thật để phát Bước 5
config/ Chưa có gì cần cấu hình ngoài DSN Bước 4
logger/ log chuẩn là đủ cho tới khi có HTTP handler Bước 4
httpx/ Chưa có endpoint nào Bước 4

Vì sao không viết sẵn hết một lượt cho xong?

Code viết trước khi có người dùng nó là code chưa được kiểm chứng. Bạn sẽ đoán sai chữ ký hàm, đoán sai thứ cần trừu tượng, rồi viết lại — nhưng lần này phải mang theo cả những chỗ đã lỡ dùng sai.

Ba package ở bước này được viết ngay bây giờ vì mọi thứ sau đó đều cần chúng. httpx thì không — viết nó lúc chưa có handler nào là thiết kế trong bóng tối.


11. Checklist hoàn thành

  • [ ] go build ./... không lỗi
  • [ ] go vet ./... sạch
  • [ ] go test ./internal/... -race -count=1 — tất cả ok, không SKIP
  • [ ] go run ./cmd/smoketest in đủ 3 dấu ✅
  • [ ] contracts không import package nội bộ nào (kiểm tra bên dưới)
  • [ ] platform/ không import modules/ (chưa có modules/, nhưng nhớ luật)

Kiểm tra contracts đúng là lá của đồ thị phụ thuộc:

go list -deps ./internal/contracts/... | Select-String "yourname/community"

Chỉ được ra chính các package contracts. Nếu có bất kỳ dòng nào khác — đặc biệt là platform hay modules — bạn đã tạo ra một phụ thuộc sẽ thành vòng lặp ở bước sau.

Cấu trúc thư mục lúc này:

project/
├── cmd/smoketest/main.go
├── deployments/docker-compose.yml
├── docs/
│   ├── ARCHITECTURE.md
│   ├── 01-khoi-tao-du-an-va-ha-tang.md
│   └── 02-nen-mong-platform.md
├── internal/
│   ├── contracts/
│   │   ├── envelope.go
│   │   ├── envelope_test.go
│   │   ├── topics.go
│   │   └── v1/
│   │       ├── doc.go
│   │       └── post_events.go
│   └── platform/
│       ├── eventbus/
│       │   ├── bus.go
│       │   ├── errors.go
│       │   ├── router.go
│       │   └── inmem/
│       │       ├── bus.go
│       │       └── bus_test.go
│       └── postgres/
│           ├── pool.go
│           ├── errors.go
│           ├── tx.go
│           └── tx_test.go
├── go.mod
└── go.sum

Commit:

git add .
git commit -m "feat(platform): dựng nền móng postgres, contracts và eventbus"

Bước tiếp theo

Bước 3 — Nền tảng dữ liệu và Outbox: migration cho users, user_stats, posts, outbox, consumed_events; cấu hình sqlc; viết platform/outbox (writer + relay) và platform/idempotency dựa trên TxManager vừa xong.

Đó là bước biến Luật 3 từ nguyên tắc trên giấy thành code chạy được.


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í