477 lines
11 KiB
Markdown
477 lines
11 KiB
Markdown
|
||||
|
|
go-casstm
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
# go-caatsm Refactor Plan
|
|||
|
|
|
|||
|
|
## Objective
|
|||
|
|
|
|||
|
|
Refactor the project to adopt a modern, maintainable, and scalable architecture using:
|
|||
|
|
|
|||
|
|
- Clean Architecture (app / domain / adapter / infra)
|
|||
|
|
|
|||
|
|
- nats.go JetStream (replace Watermill)
|
|||
|
|
|
|||
|
|
- PostgreSQL pgx (replace Hasura GraphQL)
|
|||
|
|
|
|||
|
|
- Koanf configuration system (replace Viper)
|
|||
|
|
|
|||
|
|
- Google Wire for dependency injection
|
|||
|
|
|
|||
|
|
- Structured logging (+ optional metrics/tracing)
|
|||
|
|
|
|||
|
|
|
|||
|
|
Goal: improve reliability, performance, extensibility, and professional engineering quality.
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
## High-Level Architecture
|
|||
|
|
|
|||
|
|
Refactor into the following structure:
|
|||
|
|
|
|||
|
|
```text
|
|||
|
|
/cmd/receiver/main.go # entrypoint using wire-generated injector
|
|||
|
|
/config/config.toml
|
|||
|
|
/internal
|
|||
|
|
/app # Orchestrates flows
|
|||
|
|
processor.go
|
|||
|
|
listener.go
|
|||
|
|
/domain
|
|||
|
|
telegram.go
|
|||
|
|
/adapter
|
|||
|
|
parser/
|
|||
|
|
mapper/
|
|||
|
|
/infra
|
|||
|
|
config/ # koanf loader
|
|||
|
|
nats/ # jetstream consumer/publisher
|
|||
|
|
postgres/ # pgx repository
|
|||
|
|
log/ # zap logger
|
|||
|
|
/pkg/di/wire.go # wire DI root
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
Principles:
|
|||
|
|
|
|||
|
|
- Domain is pure Go types (no external imports).
|
|||
|
|
|
|||
|
|
- App orchestrates: NATS msg → parser → domain → repository.
|
|||
|
|
|
|||
|
|
- Infra handles external concerns (NATS, PostgreSQL, config, logging).
|
|||
|
|
|
|||
|
|
- Adapter performs mapping between infra/domain.
|
|||
|
|
|
|||
|
|
- `cmd` 只负责启动,不包含业务逻辑。
|
|||
|
|
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
## Phase 1 — Project Structure Migration
|
|||
|
|
|
|||
|
|
**Goal:** Introduce new directories without breaking existing code.
|
|||
|
|
|
|||
|
|
### Tasks
|
|||
|
|
|
|||
|
|
- Create new `/internal/app`, `/internal/domain`, `/internal/adapter`, `/internal/infra` directories.
|
|||
|
|
|
|||
|
|
- Move domain-level structs (telegram, metadata) into `/internal/domain`.
|
|||
|
|
|
|||
|
|
- Move parsing logic into `/internal/adapter/parser`.
|
|||
|
|
|
|||
|
|
- Add `/pkg/di` for Wire.
|
|||
|
|
|
|||
|
|
- Update `go.mod` and imports accordingly.
|
|||
|
|
|
|||
|
|
|
|||
|
|
### Acceptance Criteria
|
|||
|
|
|
|||
|
|
- Project builds successfully.
|
|||
|
|
|
|||
|
|
- Existing behavior unchanged(只是结构调整,不改逻辑).
|
|||
|
|
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
## Phase 2 — Replace Viper → Koanf
|
|||
|
|
|
|||
|
|
**Goal:** Introduce reliable & explicit config loading.
|
|||
|
|
|
|||
|
|
### Tasks
|
|||
|
|
|
|||
|
|
- Add Koanf loader at `/internal/infra/config/koanf.go`.
|
|||
|
|
|
|||
|
|
- Load from file (`config/config.toml`) then environment (`CAATSM_` prefix).
|
|||
|
|
|
|||
|
|
- Define a strongly typed `Config` struct (NATS, Postgres, logging, etc.).
|
|||
|
|
|
|||
|
|
- Remove global singleton config; pass `*Config` explicitly via DI.
|
|||
|
|
|
|||
|
|
- Add config validation logic (e.g. non-empty URLs, timeouts > 0).
|
|||
|
|
|
|||
|
|
|
|||
|
|
### Example (参考实现思路)
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
func LoadConfig() (*Config, error) {
|
|||
|
|
k := koanf.New(".")
|
|||
|
|
|
|||
|
|
if err := k.Load(file.Provider("config/config.toml"), toml.Parser()); err != nil {
|
|||
|
|
return nil, err
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
if err := k.Load(env.Provider("CAATSM_", ".", func(s string) string {
|
|||
|
|
return strings.ToLower(strings.TrimPrefix(s, "CAATSM_"))
|
|||
|
|
}), nil); err != nil {
|
|||
|
|
return nil, err
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
var cfg Config
|
|||
|
|
if err := k.Unmarshal("", &cfg); err != nil {
|
|||
|
|
return nil, err
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
return &cfg, cfg.Validate()
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### Acceptance Criteria
|
|||
|
|
|
|||
|
|
- Running `go run cmd/receiver/main.go` loads config via Koanf correctly。
|
|||
|
|
|
|||
|
|
- No global config singletons remain。
|
|||
|
|
|
|||
|
|
- Unit tests can construct `Config` directly,方便单测。
|
|||
|
|
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
## Phase 3 — Wire Dependency Injection
|
|||
|
|
|
|||
|
|
**Goal:** Remove manual wiring logic, centralize dependency creation.
|
|||
|
|
|
|||
|
|
### Tasks
|
|||
|
|
|
|||
|
|
- Create `/pkg/di/wire.go` with injectors.
|
|||
|
|
|
|||
|
|
- Provide constructors:
|
|||
|
|
|
|||
|
|
- `ProvideConfig` (Koanf)
|
|||
|
|
|
|||
|
|
- `ProvideLogger` (Zap)
|
|||
|
|
|
|||
|
|
- `ProvideJetStream` (NATS)
|
|||
|
|
|
|||
|
|
- `ProvideDB` (pgxpool)
|
|||
|
|
|
|||
|
|
- `ProvideRepository` (Postgres repo)
|
|||
|
|
|
|||
|
|
- `NewMessageProcessor` (app layer)
|
|||
|
|
|
|||
|
|
- Generate `wire_gen.go`.
|
|||
|
|
|
|||
|
|
- Modify `cmd/receiver/main.go` to use Wire-generated `Initialize()` (或类似函数名)。
|
|||
|
|
|
|||
|
|
|
|||
|
|
### Example Wire skeleton
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
//go:build wireinject
|
|||
|
|
|
|||
|
|
package di
|
|||
|
|
|
|||
|
|
import (
|
|||
|
|
"github.com/google/wire"
|
|||
|
|
"go-caatsm/internal/app"
|
|||
|
|
"go-caatsm/internal/infra/config"
|
|||
|
|
"go-caatsm/internal/infra/log"
|
|||
|
|
"go-caatsm/internal/infra/nats"
|
|||
|
|
"go-caatsm/internal/infra/postgres"
|
|||
|
|
)
|
|||
|
|
|
|||
|
|
func InitializeProcessor() (*app.MessageProcessor, error) {
|
|||
|
|
wire.Build(
|
|||
|
|
config.ProvideConfig,
|
|||
|
|
log.ProvideLogger,
|
|||
|
|
nats.ProvideJetStream,
|
|||
|
|
postgres.ProvideDB,
|
|||
|
|
postgres.ProvideRepository,
|
|||
|
|
app.NewMessageProcessor,
|
|||
|
|
)
|
|||
|
|
return &app.MessageProcessor{}, nil
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### Acceptance Criteria
|
|||
|
|
|
|||
|
|
- Project builds with Wire DI。
|
|||
|
|
|
|||
|
|
- main.go 只负责调用 `InitializeProcessor()` 和启动 processor。
|
|||
|
|
|
|||
|
|
- 新增依赖时只需修改 Wire graph,不用手动改 main.go。
|
|||
|
|
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
## Phase 4 — Replace Watermill → nats.go JetStream
|
|||
|
|
|
|||
|
|
**Goal:** Gain full control over message flow, retries, DLQ.
|
|||
|
|
|
|||
|
|
### Tasks
|
|||
|
|
|
|||
|
|
- 引入 `/internal/infra/nats/jetstream.go`,实现:
|
|||
|
|
|
|||
|
|
- 连接创建(`nats.Connect`,`js, _ := nc.JetStream()`)
|
|||
|
|
|
|||
|
|
- Stream + Consumer 自动创建(如不存在则创建)
|
|||
|
|
|
|||
|
|
- 使用 Pull Subscribe 模式(`PullSubscribe`)
|
|||
|
|
|
|||
|
|
- 手动 ACK / NAK
|
|||
|
|
|
|||
|
|
- 简单 Retry 策略(MaxDeliveries + NAK)
|
|||
|
|
|
|||
|
|
- 死信队列(DLQ stream/subject)
|
|||
|
|
|
|||
|
|
- 实现批量抓取(例如 `Fetch(50, MaxWait(...))`)。
|
|||
|
|
|
|||
|
|
- 实现 `Consumer.Start(ctx)`,内部循环读取消息并调用 `app.MessageProcessor.Handle()`。
|
|||
|
|
|
|||
|
|
|
|||
|
|
### Example 消费逻辑骨架
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
func (c *Consumer) Start(ctx context.Context) error {
|
|||
|
|
sub, err := c.js.PullSubscribe(c.subject, c.consumerName)
|
|||
|
|
if err != nil {
|
|||
|
|
return err
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
for {
|
|||
|
|
select {
|
|||
|
|
case <-ctx.Done():
|
|||
|
|
return ctx.Err()
|
|||
|
|
default:
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
msgs, err := sub.Fetch(50, nats.MaxWait(2*time.Second))
|
|||
|
|
if err != nil {
|
|||
|
|
if errors.Is(err, nats.ErrTimeout) {
|
|||
|
|
continue
|
|||
|
|
}
|
|||
|
|
c.logger.Error("fetch failed", zap.Error(err))
|
|||
|
|
continue
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
for _, msg := range msgs {
|
|||
|
|
if err := c.handler.Handle(ctx, msg.Data); err != nil {
|
|||
|
|
_ = msg.Nak()
|
|||
|
|
continue
|
|||
|
|
}
|
|||
|
|
_ = msg.Ack()
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### Acceptance Criteria
|
|||
|
|
|
|||
|
|
- 消费逻辑完全基于 nats.go,不再依赖 Watermill。
|
|||
|
|
|
|||
|
|
- ACK / NAK 正常工作,可通过 JetStream 管理界面/CLI 查看重试与 DLQ。
|
|||
|
|
|
|||
|
|
- 可通过配置控制批量大小、等待时间、MaxDeliveries 等。
|
|||
|
|
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
## Phase 5 — Replace Hasura GraphQL → PostgreSQL (pgx)
|
|||
|
|
|
|||
|
|
**Goal:** High-performance & reliable write pipeline.
|
|||
|
|
|
|||
|
|
### Tasks
|
|||
|
|
|
|||
|
|
- 添加 `/internal/infra/postgres/db.go`,使用 `pgxpool.Pool` 管理连接。
|
|||
|
|
|
|||
|
|
- 添加 `/internal/infra/postgres/repository.go`:
|
|||
|
|
|
|||
|
|
- `InsertOne(ctx, telegram domain.Telegram) error`
|
|||
|
|
|
|||
|
|
- `InsertBatch(ctx, []domain.Telegram) error`(使用 `CopyFrom`)
|
|||
|
|
|
|||
|
|
- 定义 telegram 表结构(如已存在则对齐 struct 和列)。
|
|||
|
|
|
|||
|
|
- 增加必要索引(如 `uuid`、时间戳、业务 key 等)。
|
|||
|
|
|
|||
|
|
- 删除 Hasura GraphQL client、genqlient 相关代码。
|
|||
|
|
|
|||
|
|
|
|||
|
|
### Example CopyFrom 骨架
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
func (r *Repository) InsertBatch(ctx context.Context, msgs []domain.Telegram) error {
|
|||
|
|
rows := make([][]any, len(msgs))
|
|||
|
|
for i, m := range msgs {
|
|||
|
|
rows[i] = []any{
|
|||
|
|
m.UUID,
|
|||
|
|
m.Raw,
|
|||
|
|
m.ParsedJSON,
|
|||
|
|
m.CreatedAt,
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
_, err := r.pool.CopyFrom(
|
|||
|
|
ctx,
|
|||
|
|
pgx.Identifier{"aviation_telegrams"},
|
|||
|
|
[]string{"uuid", "raw", "parsed", "created_at"},
|
|||
|
|
pgx.CopyFromRows(rows),
|
|||
|
|
)
|
|||
|
|
return err
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### Acceptance Criteria
|
|||
|
|
|
|||
|
|
- 消息数据成功写入 PostgreSQL。
|
|||
|
|
|
|||
|
|
- 批量写入时使用 CopyFrom,性能明显优于单条 INSERT。
|
|||
|
|
|
|||
|
|
- Hasura / GraphQL 相关依赖从代码和 go.mod 中移除。
|
|||
|
|
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
## Phase 6 — Application Layer (Processor)
|
|||
|
|
|
|||
|
|
**Goal:** Create clean orchestrator for the message lifecycle.
|
|||
|
|
|
|||
|
|
### Tasks
|
|||
|
|
|
|||
|
|
- 在 `/internal/app/processor.go` 实现 `MessageProcessor`:
|
|||
|
|
|
|||
|
|
- 接口定义:
|
|||
|
|
|
|||
|
|
- `type Parser interface { Parse(raw []byte) (domain.Telegram, error) }`
|
|||
|
|
|
|||
|
|
- `type Repository interface { InsertOne / InsertBatch }`
|
|||
|
|
|
|||
|
|
- 核心流程:
|
|||
|
|
|
|||
|
|
1. 收到 NATS 消息(由 consumer 调用 `HandleMessage` 或类似接口)
|
|||
|
|
|
|||
|
|
2. 调用 `Parser.Parse` 得到 `domain.Telegram`
|
|||
|
|
|
|||
|
|
3. 调用 `Repository.Insert...` 写入数据库
|
|||
|
|
|
|||
|
|
4. 返回成功/失败,由 caller 决定 ACK/NAK
|
|||
|
|
|
|||
|
|
- 在 `/internal/adapter/parser` 中处理具体报文解析逻辑,保持 domain 纯净。
|
|||
|
|
|
|||
|
|
|
|||
|
|
### Acceptance Criteria
|
|||
|
|
|
|||
|
|
- Processor 不依赖具体的 NATS / pgx 类型,只依赖接口。
|
|||
|
|
|
|||
|
|
- Parser / Repository 可以在测试中替换为 mock。
|
|||
|
|
|
|||
|
|
- 业务流程清晰、单一职责。
|
|||
|
|
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
## Phase 7 — Logging & Observability
|
|||
|
|
|
|||
|
|
**Goal:** Unify logging and enable production-ready debugging.
|
|||
|
|
|
|||
|
|
### Tasks
|
|||
|
|
|
|||
|
|
- 在 `/internal/infra/log/logger.go` 实现 Zap 初始化(支持 dev/prod 模式)。
|
|||
|
|
|
|||
|
|
- 将 main、consumer、processor、repository 中的 `fmt.Println` 替换为结构化日志。
|
|||
|
|
|
|||
|
|
- 每条关键日志附加必要 context 字段:
|
|||
|
|
|
|||
|
|
- `message_id`
|
|||
|
|
|
|||
|
|
- `subject`
|
|||
|
|
|
|||
|
|
- `stream`
|
|||
|
|
|
|||
|
|
- `attempt`
|
|||
|
|
|
|||
|
|
- (可选)添加 Prometheus metrics(处理量、错误数、重试次数)。
|
|||
|
|
|
|||
|
|
|
|||
|
|
### Acceptance Criteria
|
|||
|
|
|
|||
|
|
- 日志输出统一,方便在 Loki / ELK 中检索。
|
|||
|
|
|
|||
|
|
- 出错时能通过日志快速定位是哪个环节(NATS 消费 / 解析 / DB 写入)出了问题。
|
|||
|
|
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
## Phase 8 — Remove Dead Code & Cleanup
|
|||
|
|
|
|||
|
|
**Goal:** Remove legacy patterns and unused modules.
|
|||
|
|
|
|||
|
|
### Tasks
|
|||
|
|
|
|||
|
|
- 移除 Watermill 相关代码与依赖。
|
|||
|
|
|
|||
|
|
- 移除 Hasura / genqlient 相关代码与依赖。
|
|||
|
|
|
|||
|
|
- 移除 Viper 配置加载器与全局单例。
|
|||
|
|
|
|||
|
|
- 删除不再使用的 handler / repository 实现。
|
|||
|
|
|
|||
|
|
- 运行 `go mod tidy` 清理依赖。
|
|||
|
|
|
|||
|
|
- 检查 Taskfile / Makefile,更新为新的启动、测试命令。
|
|||
|
|
|
|||
|
|
|
|||
|
|
### Acceptance Criteria
|
|||
|
|
|
|||
|
|
- `go test ./...` 与 `go build ./...` 均成功。
|
|||
|
|
|
|||
|
|
- go.mod 中不再包含 Watermill / Hasura / genqlient / Viper。
|
|||
|
|
|
|||
|
|
- 代码中不再有全局 Config/Logger 单例。
|
|||
|
|
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
## Final Acceptance Criteria
|
|||
|
|
|
|||
|
|
Refactor 完成的标志:
|
|||
|
|
|
|||
|
|
1. **启动链路:**
|
|||
|
|
|
|||
|
|
- 使用 Koanf 加载配置。
|
|||
|
|
|
|||
|
|
- 使用 Wire 完成依赖注入。
|
|||
|
|
|
|||
|
|
- 使用 nats.go JetStream 消费消息。
|
|||
|
|
|
|||
|
|
- 使用 pgx 将数据写入 PostgreSQL。
|
|||
|
|
|
|||
|
|
2. **架构层次清晰:**
|
|||
|
|
|
|||
|
|
- `internal/domain` 无外部依赖。
|
|||
|
|
|
|||
|
|
- `internal/app` 只依赖 domain + 抽象接口。
|
|||
|
|
|
|||
|
|
- `internal/infra` 只负责技术细节。
|
|||
|
|
|
|||
|
|
- `cmd` 只启动,不包含业务逻辑。
|
|||
|
|
|
|||
|
|
3. **旧技术栈完全移除:**
|
|||
|
|
|
|||
|
|
- Watermill、Hasura、GraphQL、Viper、全局单例全部删除。
|
|||
|
|
|
|||
|
|
4. **数据流全链路可工作:**
|
|||
|
|
|
|||
|
|
- NATS → Parser → Domain Model → Repository → PostgreSQL 全流程可验证。
|
|||
|
|
|
|||
|
|
|
|||
|
|
---
|
|||
|
|
|