commit 976ec3ca70bb807c073a83b8569e76343bee54b1 Author: ryan Date: Sat Jun 27 12:57:11 2026 +0800 feat: project scaffold, config loading, skeleton main diff --git a/.claude/settings.local.json b/.claude/settings.local.json new file mode 100644 index 0000000..363efcc --- /dev/null +++ b/.claude/settings.local.json @@ -0,0 +1,11 @@ +{ + "permissions": { + "allow": [ + "Bash(go mod *)", + "Bash(go build *)", + "Bash(go run *)", + "Bash(git add *)", + "Bash(git commit *)" + ] + } +} diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..cd63957 --- /dev/null +++ b/.gitignore @@ -0,0 +1,4 @@ +.superpowers/ +bin/ +*.log +.DS_Store diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..ca20c1a --- /dev/null +++ b/Makefile @@ -0,0 +1,25 @@ +.PHONY: build run test docker-build docker-up docker-down migrate-up migrate-down + +build: + go build -o bin/server ./cmd/server + +run: + go run ./cmd/server + +test: + go test ./internal/... -v -count=1 + +docker-build: + docker build -t notification-service . + +docker-up: + docker-compose up -d + +docker-down: + docker-compose down + +migrate-up: + migrate -path migrations -database "mysql://notify:notify@tcp(127.0.0.1:3306)/notification" up + +migrate-down: + migrate -path migrations -database "mysql://notify:notify@tcp(127.0.0.1:3306)/notification" down diff --git a/cmd/server/main.go b/cmd/server/main.go new file mode 100644 index 0000000..127c512 --- /dev/null +++ b/cmd/server/main.go @@ -0,0 +1,21 @@ +package main + +import ( + "log/slog" + "os" + + "aiaa-notification-service/internal/config" +) + +func main() { + logger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})) + slog.SetDefault(logger) + + cfg, err := config.Load("config/config.yaml") + if err != nil { + slog.Error("failed to load config", "error", err) + os.Exit(1) + } + + slog.Info("config loaded", "port", cfg.Server.Port) +} diff --git a/config/config.yaml b/config/config.yaml new file mode 100644 index 0000000..c7d3905 --- /dev/null +++ b/config/config.yaml @@ -0,0 +1,26 @@ +server: + port: 8080 + admin_key: "admin-sk-change-me" + +database: + host: "127.0.0.1" + port: 3306 + user: "notify" + password: "${DB_PASSWORD:-notify}" + database: "notification" + +redis: + host: "127.0.0.1" + port: 6379 + password: "" + db: 0 + +smtp: + host: "smtp.example.com" + port: 587 + user: "notify@example.com" + password: "${SMTP_PASSWORD:-}" + from: "Notification Service " + +rate_limit: + default: 100 diff --git a/docs/superpowers/plans/2025-06-27-notification-service-plan.md b/docs/superpowers/plans/2025-06-27-notification-service-plan.md new file mode 100644 index 0000000..ba71d8d --- /dev/null +++ b/docs/superpowers/plans/2025-06-27-notification-service-plan.md @@ -0,0 +1,3724 @@ +# Notification Service Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Build a pure-Go notification service — webhook in, parse (JSON/Regex/Text), match rules, evaluate conditions, render Go templates, dispatch to DingTalk/WeCom/Email/Bark. + +**Architecture:** Gin HTTP → auth+ratelimit middleware → handler → parser → engine(matcher→condition→renderer→router) → adapter(ChannelSender interface) → external services. MySQL for config, Redis for cache+ratelimit. + +**Tech Stack:** Go 1.22+, gin, sqlx, go-redis/v9, text/template, viper, slog, golang-migrate, MySQL 8.0, Redis 7 + +## Global Constraints + +- Go 1.22+, module path: `aiaa-notification-service` +- MySQL 8.0 InnoDB, Redis 7 +- Config via `config/config.yaml` with `${ENV_VAR}` interpolation +- Tests: `testing` stdlib + `httptest` for handlers, `testify/assert` for assertions +- Docker multi-stage: `golang:1.22-alpine` → `alpine:3.20` + +## File Structure + +``` +notification-service/ +├── cmd/server/main.go +├── config/config.yaml +├── migrations/ +│ ├── 001_init.up.sql +│ └── 001_init.down.sql +├── internal/ +│ ├── config/config.go +│ ├── model/model.go +│ ├── store/ +│ │ ├── mysql.go +│ │ ├── source.go +│ │ ├── template.go +│ │ ├── channel.go +│ │ ├── rule.go +│ │ └── message_log.go +│ ├── cache/redis.go +│ ├── parser/ +│ │ ├── parser.go ← Parser interface + factory +│ │ ├── json.go +│ │ ├── regex.go +│ │ └── text.go +│ ├── condition/ +│ │ └── evaluator.go ← condition array evaluator +│ ├── adapter/ +│ │ ├── adapter.go +│ │ ├── dingtalk.go +│ │ ├── wecom.go +│ │ ├── email.go +│ │ └── bark.go +│ ├── engine/ +│ │ ├── matcher.go +│ │ ├── renderer.go +│ │ └── router.go +│ ├── handler/ +│ │ ├── middleware.go +│ │ ├── notify.go +│ │ ├── source.go +│ │ ├── template.go +│ │ ├── channel.go +│ │ ├── rule.go +│ │ └── message_log.go +│ └── retry/retry.go +├── Dockerfile +├── docker-compose.yml +├── Makefile +└── go.mod +``` + +--- + +### Task 1: Project Scaffold & Configuration Loading + +**Files:** +- Create: `go.mod` +- Create: `cmd/server/main.go` (skeleton) +- Create: `config/config.yaml` +- Create: `internal/config/config.go` +- Create: `Makefile` + +**Interfaces:** +- Produces: `config.Config` struct, `config.Load(path string) (*Config, error)` + +- [ ] **Step 1: Initialize Go module** + +```bash +cd /Users/ryan/Documents/code/go/aiaa-notification-service +go mod init aiaa-notification-service +``` + +- [ ] **Step 2: Create config file** + +File: `config/config.yaml` +```yaml +server: + port: 8080 + admin_key: "admin-sk-change-me" + +database: + host: "127.0.0.1" + port: 3306 + user: "notify" + password: "${DB_PASSWORD:-notify}" + database: "notification" + +redis: + host: "127.0.0.1" + port: 6379 + password: "" + db: 0 + +smtp: + host: "smtp.example.com" + port: 587 + user: "notify@example.com" + password: "${SMTP_PASSWORD:-}" + from: "Notification Service " + +rate_limit: + default: 100 +``` + +- [ ] **Step 3: Write Config loader** + +File: `internal/config/config.go` +```go +package config + +import ( + "fmt" + "os" + "strings" + + "github.com/spf13/viper" +) + +type Config struct { + Server ServerConfig `mapstructure:"server"` + Database DatabaseConfig `mapstructure:"database"` + Redis RedisConfig `mapstructure:"redis"` + SMTP SMTPConfig `mapstructure:"smtp"` + RateLimit RateLimitConfig `mapstructure:"rate_limit"` +} + +type ServerConfig struct { + Port int `mapstructure:"port"` + AdminKey string `mapstructure:"admin_key"` +} + +type DatabaseConfig struct { + Host string `mapstructure:"host"` + Port int `mapstructure:"port"` + User string `mapstructure:"user"` + Password string `mapstructure:"password"` + Database string `mapstructure:"database"` +} + +func (d DatabaseConfig) DSN() string { + return fmt.Sprintf("%s:%s@tcp(%s:%d)/%s?charset=utf8mb4&parseTime=true&loc=Local", + d.User, d.Password, d.Host, d.Port, d.Database) +} + +type RedisConfig struct { + Host string `mapstructure:"host"` + Port int `mapstructure:"port"` + Password string `mapstructure:"password"` + DB int `mapstructure:"db"` +} + +func (r RedisConfig) Addr() string { + return fmt.Sprintf("%s:%d", r.Host, r.Port) +} + +type SMTPConfig struct { + Host string `mapstructure:"host"` + Port int `mapstructure:"port"` + User string `mapstructure:"user"` + Password string `mapstructure:"password"` + From string `mapstructure:"from"` +} + +type RateLimitConfig struct { + Default int `mapstructure:"default"` +} + +func Load(path string) (*Config, error) { + v := viper.New() + v.SetConfigFile(path) + v.SetEnvPrefix("NOTIFY") + v.SetEnvKeyReplacer(strings.NewReplacer(".", "_")) + v.AutomaticEnv() + + if err := v.ReadInConfig(); err != nil { + return nil, fmt.Errorf("read config: %w", err) + } + + // Expand ${ENV_VAR} placeholders in config values + for _, key := range v.AllKeys() { + val := v.GetString(key) + if strings.Contains(val, "${") { + expanded := expandEnv(val) + v.Set(key, expanded) + } + } + + var cfg Config + if err := v.Unmarshal(&cfg); err != nil { + return nil, fmt.Errorf("unmarshal config: %w", err) + } + return &cfg, nil +} + +func expandEnv(s string) string { + return os.Expand(s, func(key string) string { + // support ${VAR:-default} + if i := strings.Index(key, ":-"); i >= 0 { + name := key[:i] + def := key[i+2:] + if v, ok := os.LookupEnv(name); ok { + return v + } + return def + } + return os.Getenv(key) + }) +} +``` + +- [ ] **Step 4: Create skeleton main.go** + +File: `cmd/server/main.go` +```go +package main + +import ( + "log/slog" + "os" + + "aiaa-notification-service/internal/config" +) + +func main() { + logger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})) + slog.SetDefault(logger) + + cfg, err := config.Load("config/config.yaml") + if err != nil { + slog.Error("failed to load config", "error", err) + os.Exit(1) + } + + slog.Info("config loaded", "port", cfg.Server.Port) +} +``` + +- [ ] **Step 5: Create Makefile** + +File: `Makefile` +```makefile +.PHONY: build run test docker-build docker-up docker-down migrate-up migrate-down + +build: + go build -o bin/server ./cmd/server + +run: + go run ./cmd/server + +test: + go test ./internal/... -v -count=1 + +docker-build: + docker build -t notification-service . + +docker-up: + docker-compose up -d + +docker-down: + docker-compose down + +migrate-up: + migrate -path migrations -database "mysql://notify:notify@tcp(127.0.0.1:3306)/notification" up + +migrate-down: + migrate -path migrations -database "mysql://notify:notify@tcp(127.0.0.1:3306)/notification" down +``` + +- [ ] **Step 6: Install dependencies and verify build** + +```bash +go mod tidy +go build ./... +``` +Expected: build succeeds with no errors. + +- [ ] **Step 7: Commit** + +```bash +git add -A && git commit -m "feat: project scaffold, config loading, skeleton main" +``` + +--- + +### Task 2: Models & Database Migration + +**Files:** +- Create: `internal/model/model.go` +- Create: `migrations/001_init.up.sql` +- Create: `migrations/001_init.down.sql` + +**Interfaces:** +- Produces: `model.Source`, `model.Template`, `model.Channel`, `model.Rule`, `model.RuleChannel`, `model.MessageLog`, `model.Condition` + +- [ ] **Step 1: Write Go models** + +File: `internal/model/model.go` +```go +package model + +import ( + "database/sql" + "encoding/json" + "time" +) + +type Source struct { + ID int `db:"id" json:"id"` + Name string `db:"name" json:"name"` + APIKey string `db:"api_key" json:"api_key,omitempty"` + ParseMode string `db:"parse_mode" json:"parse_mode"` + ParsePattern string `db:"parse_pattern" json:"parse_pattern,omitempty"` + Status int `db:"status" json:"status"` + CreatedAt time.Time `db:"created_at" json:"created_at"` + UpdatedAt time.Time `db:"updated_at" json:"updated_at"` +} + +type Template struct { + ID int `db:"id" json:"id"` + Name string `db:"name" json:"name"` + Content string `db:"content" json:"content"` + CreatedAt time.Time `db:"created_at" json:"created_at"` + UpdatedAt time.Time `db:"updated_at" json:"updated_at"` +} + +type Channel struct { + ID int `db:"id" json:"id"` + Name string `db:"name" json:"name"` + Type string `db:"type" json:"type"` + Config *json.RawMessage `db:"config" json:"config"` + Status int `db:"status" json:"status"` + CreatedAt time.Time `db:"created_at" json:"created_at"` + UpdatedAt time.Time `db:"updated_at" json:"updated_at"` +} + +type Rule struct { + ID int `db:"id" json:"id"` + SourceID int `db:"source_id" json:"source_id"` + Event string `db:"event" json:"event"` + TemplateID int `db:"template_id" json:"template_id"` + Conditions *json.RawMessage `db:"conditions" json:"conditions,omitempty"` + Enabled int `db:"enabled" json:"enabled"` + CreatedAt time.Time `db:"created_at" json:"created_at"` + UpdatedAt time.Time `db:"updated_at" json:"updated_at"` +} + +type Condition struct { + Field string `json:"field"` + Op string `json:"op"` + Value string `json:"value,omitempty"` +} + +type RuleChannel struct { + ID int `db:"id" json:"id"` + RuleID int `db:"rule_id" json:"rule_id"` + ChannelID int `db:"channel_id" json:"channel_id"` + Enabled int `db:"enabled" json:"enabled"` +} + +type MessageLog struct { + ID int64 `db:"id" json:"id"` + RuleID int `db:"rule_id" json:"rule_id"` + ChannelID int `db:"channel_id" json:"channel_id"` + Source string `db:"source" json:"source"` + Event string `db:"event" json:"event"` + Payload json.RawMessage `db:"payload" json:"payload"` + Content string `db:"content" json:"content"` + Status string `db:"status" json:"status"` + RetryCount int `db:"retry_count" json:"retry_count"` + Response sql.NullString `db:"response" json:"response,omitempty"` + ErrorMsg sql.NullString `db:"error_msg" json:"error_msg,omitempty"` + CreatedAt time.Time `db:"created_at" json:"created_at"` +} +``` + +- [ ] **Step 2: Write migration SQL** + +File: `migrations/001_init.up.sql` +```sql +CREATE TABLE source ( + id INT AUTO_INCREMENT PRIMARY KEY, + name VARCHAR(64) NOT NULL, + api_key VARCHAR(128) NOT NULL, + parse_mode VARCHAR(16) NOT NULL DEFAULT 'json', + parse_pattern VARCHAR(512) DEFAULT NULL, + status TINYINT NOT NULL DEFAULT 1, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, + UNIQUE KEY uk_name (name), + UNIQUE KEY uk_api_key (api_key) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE template ( + id INT AUTO_INCREMENT PRIMARY KEY, + name VARCHAR(64) NOT NULL, + content TEXT NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, + UNIQUE KEY uk_name (name) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE channel ( + id INT AUTO_INCREMENT PRIMARY KEY, + name VARCHAR(32) NOT NULL, + type VARCHAR(32) NOT NULL, + config JSON NOT NULL, + status TINYINT NOT NULL DEFAULT 1, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, + UNIQUE KEY uk_name (name) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE rule ( + id INT AUTO_INCREMENT PRIMARY KEY, + source_id INT NOT NULL, + event VARCHAR(64) NOT NULL, + template_id INT NOT NULL, + conditions JSON DEFAULT NULL, + enabled TINYINT NOT NULL DEFAULT 1, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, + UNIQUE KEY uk_source_event (source_id, event), + FOREIGN KEY (source_id) REFERENCES source(id), + FOREIGN KEY (template_id) REFERENCES template(id) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE rule_channel ( + id INT AUTO_INCREMENT PRIMARY KEY, + rule_id INT NOT NULL, + channel_id INT NOT NULL, + enabled TINYINT NOT NULL DEFAULT 1, + UNIQUE KEY uk_rule_channel (rule_id, channel_id), + FOREIGN KEY (rule_id) REFERENCES rule(id) ON DELETE CASCADE, + FOREIGN KEY (channel_id) REFERENCES channel(id) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE message_log ( + id BIGINT AUTO_INCREMENT PRIMARY KEY, + rule_id INT NOT NULL, + channel_id INT NOT NULL, + source VARCHAR(64) NOT NULL, + event VARCHAR(64) NOT NULL, + payload JSON NOT NULL, + content TEXT NOT NULL, + status ENUM('pending','success','failed','retrying') NOT NULL DEFAULT 'pending', + retry_count INT NOT NULL DEFAULT 0, + response TEXT DEFAULT NULL, + error_msg TEXT DEFAULT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + INDEX idx_source_event_time (source, event, created_at) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; +``` + +File: `migrations/001_init.down.sql` +```sql +DROP TABLE IF EXISTS message_log; +DROP TABLE IF EXISTS rule_channel; +DROP TABLE IF EXISTS rule; +DROP TABLE IF EXISTS channel; +DROP TABLE IF EXISTS template; +DROP TABLE IF EXISTS source; +``` + +- [ ] **Step 3: Verify models compile** + +```bash +go build ./internal/model/... +``` +Expected: no errors. + +- [ ] **Step 4: Commit** + +```bash +git add -A && git commit -m "feat: data models and database migration" +``` + +--- + +### Task 3: MySQL Store — Connection & Source CRUD + +**Files:** +- Create: `internal/store/mysql.go` +- Create: `internal/store/source.go` + +**Interfaces:** +- Consumes: `config.Config`, `model.Source` +- Produces: `store.NewStore(cfg *config.DatabaseConfig) (*Store, error)`, `Store.CreateSource`, `Store.GetSource`, `Store.GetSourceByAPIKey`, `Store.GetSourceByName`, `Store.ListSources`, `Store.UpdateSource`, `Store.DeleteSource` + +- [ ] **Step 1: Write store connection** + +File: `internal/store/mysql.go` +```go +package store + +import ( + "fmt" + + "aiaa-notification-service/internal/config" + + "github.com/jmoiron/sqlx" + _ "github.com/go-sql-driver/mysql" +) + +type Store struct { + DB *sqlx.DB +} + +func NewStore(cfg config.DatabaseConfig) (*Store, error) { + db, err := sqlx.Connect("mysql", cfg.DSN()) + if err != nil { + return nil, fmt.Errorf("connect mysql: %w", err) + } + db.SetMaxOpenConns(25) + db.SetMaxIdleConns(5) + return &Store{DB: db}, nil +} + +func (s *Store) Close() error { + return s.DB.Close() +} +``` + +- [ ] **Step 2: Write Source store methods** + +File: `internal/store/source.go` +```go +package store + +import ( + "context" + "crypto/rand" + "encoding/hex" + "fmt" + + "aiaa-notification-service/internal/model" +) + +func generateAPIKey() string { + b := make([]byte, 32) + rand.Read(b) + return "sk-" + hex.EncodeToString(b) +} + +func (s *Store) CreateSource(ctx context.Context, src *model.Source) error { + src.APIKey = generateAPIKey() + query := `INSERT INTO source (name, api_key, parse_mode, parse_pattern, status) VALUES (?, ?, ?, ?, ?)` + result, err := s.DB.ExecContext(ctx, query, src.Name, src.APIKey, src.ParseMode, src.ParsePattern, src.Status) + if err != nil { + return fmt.Errorf("create source: %w", err) + } + id, _ := result.LastInsertId() + src.ID = int(id) + return nil +} + +func (s *Store) GetSource(ctx context.Context, id int) (*model.Source, error) { + var src model.Source + err := s.DB.GetContext(ctx, &src, `SELECT * FROM source WHERE id = ?`, id) + if err != nil { + return nil, fmt.Errorf("get source %d: %w", id, err) + } + return &src, nil +} + +func (s *Store) GetSourceByAPIKey(ctx context.Context, apiKey string) (*model.Source, error) { + var src model.Source + err := s.DB.GetContext(ctx, &src, `SELECT * FROM source WHERE api_key = ? AND status = 1`, apiKey) + if err != nil { + return nil, fmt.Errorf("get source by api_key: %w", err) + } + return &src, nil +} + +func (s *Store) GetSourceByName(ctx context.Context, name string) (*model.Source, error) { + var src model.Source + err := s.DB.GetContext(ctx, &src, `SELECT * FROM source WHERE name = ?`, name) + if err != nil { + return nil, fmt.Errorf("get source by name %s: %w", name, err) + } + return &src, nil +} + +func (s *Store) ListSources(ctx context.Context) ([]model.Source, error) { + var sources []model.Source + err := s.DB.SelectContext(ctx, &sources, `SELECT * FROM source ORDER BY id`) + if err != nil { + return nil, fmt.Errorf("list sources: %w", err) + } + return sources, nil +} + +func (s *Store) UpdateSource(ctx context.Context, id int, src *model.Source) error { + query := `UPDATE source SET name=?, parse_mode=?, parse_pattern=?, status=? WHERE id=?` + _, err := s.DB.ExecContext(ctx, query, src.Name, src.ParseMode, src.ParsePattern, src.Status, id) + if err != nil { + return fmt.Errorf("update source %d: %w", id, err) + } + return nil +} + +func (s *Store) DeleteSource(ctx context.Context, id int) error { + _, err := s.DB.ExecContext(ctx, `DELETE FROM source WHERE id = ?`, id) + if err != nil { + return fmt.Errorf("delete source %d: %w", id, err) + } + return nil +} +``` + +- [ ] **Step 3: Install dependencies** + +```bash +go get github.com/jmoiron/sqlx github.com/go-sql-driver/mysql +go mod tidy +``` + +- [ ] **Step 4: Verify compilation** + +```bash +go build ./internal/store/... +``` +Expected: no errors. + +- [ ] **Step 5: Commit** + +```bash +git add -A && git commit -m "feat: store connection and source CRUD" +``` + +--- + +### Task 4: MySQL Store — Template & Channel CRUD + +**Files:** +- Create: `internal/store/template.go` +- Create: `internal/store/channel.go` + +**Interfaces:** +- Produces: `Store.CreateTemplate`, `Store.GetTemplate`, `Store.GetTemplateByName`, `Store.ListTemplates`, `Store.UpdateTemplate`, `Store.DeleteTemplate` +- Produces: `Store.CreateChannel`, `Store.GetChannel`, `Store.GetChannelByName`, `Store.ListChannels`, `Store.UpdateChannel`, `Store.DeleteChannel` + +- [ ] **Step 1: Write Template store** + +File: `internal/store/template.go` +```go +package store + +import ( + "context" + "fmt" + + "aiaa-notification-service/internal/model" +) + +func (s *Store) CreateTemplate(ctx context.Context, t *model.Template) error { + query := `INSERT INTO template (name, content) VALUES (?, ?)` + result, err := s.DB.ExecContext(ctx, query, t.Name, t.Content) + if err != nil { + return fmt.Errorf("create template: %w", err) + } + id, _ := result.LastInsertId() + t.ID = int(id) + return nil +} + +func (s *Store) GetTemplate(ctx context.Context, id int) (*model.Template, error) { + var t model.Template + err := s.DB.GetContext(ctx, &t, `SELECT * FROM template WHERE id = ?`, id) + if err != nil { + return nil, fmt.Errorf("get template %d: %w", id, err) + } + return &t, nil +} + +func (s *Store) GetTemplateByName(ctx context.Context, name string) (*model.Template, error) { + var t model.Template + err := s.DB.GetContext(ctx, &t, `SELECT * FROM template WHERE name = ?`, name) + if err != nil { + return nil, fmt.Errorf("get template by name %s: %w", name, err) + } + return &t, nil +} + +func (s *Store) ListTemplates(ctx context.Context) ([]model.Template, error) { + var templates []model.Template + err := s.DB.SelectContext(ctx, &templates, `SELECT * FROM template ORDER BY id`) + if err != nil { + return nil, fmt.Errorf("list templates: %w", err) + } + return templates, nil +} + +func (s *Store) UpdateTemplate(ctx context.Context, id int, t *model.Template) error { + query := `UPDATE template SET name=?, content=? WHERE id=?` + _, err := s.DB.ExecContext(ctx, query, t.Name, t.Content, id) + if err != nil { + return fmt.Errorf("update template %d: %w", id, err) + } + return nil +} + +func (s *Store) DeleteTemplate(ctx context.Context, id int) error { + _, err := s.DB.ExecContext(ctx, `DELETE FROM template WHERE id = ?`, id) + if err != nil { + return fmt.Errorf("delete template %d: %w", id, err) + } + return nil +} +``` + +- [ ] **Step 2: Write Channel store** + +File: `internal/store/channel.go` +```go +package store + +import ( + "context" + "encoding/json" + "fmt" + + "aiaa-notification-service/internal/model" +) + +func (s *Store) CreateChannel(ctx context.Context, ch *model.Channel) error { + configJSON, err := json.Marshal(ch.Config) + if err != nil { + return fmt.Errorf("marshal channel config: %w", err) + } + query := `INSERT INTO channel (name, type, config, status) VALUES (?, ?, ?, ?)` + result, err := s.DB.ExecContext(ctx, query, ch.Name, ch.Type, configJSON, ch.Status) + if err != nil { + return fmt.Errorf("create channel: %w", err) + } + id, _ := result.LastInsertId() + ch.ID = int(id) + return nil +} + +func (s *Store) GetChannel(ctx context.Context, id int) (*model.Channel, error) { + var ch model.Channel + var configBytes []byte + row := s.DB.QueryRowContext(ctx, `SELECT id, name, type, config, status, created_at, updated_at FROM channel WHERE id = ?`, id) + if err := row.Scan(&ch.ID, &ch.Name, &ch.Type, &configBytes, &ch.Status, &ch.CreatedAt, &ch.UpdatedAt); err != nil { + return nil, fmt.Errorf("get channel %d: %w", id, err) + } + raw := json.RawMessage(configBytes) + ch.Config = &raw + return &ch, nil +} + +func (s *Store) GetChannelByName(ctx context.Context, name string) (*model.Channel, error) { + var ch model.Channel + var configBytes []byte + row := s.DB.QueryRowContext(ctx, `SELECT id, name, type, config, status, created_at, updated_at FROM channel WHERE name = ?`, name) + if err := row.Scan(&ch.ID, &ch.Name, &ch.Type, &configBytes, &ch.Status, &ch.CreatedAt, &ch.UpdatedAt); err != nil { + return nil, fmt.Errorf("get channel by name %s: %w", name, err) + } + raw := json.RawMessage(configBytes) + ch.Config = &raw + return &ch, nil +} + +func (s *Store) ListChannels(ctx context.Context) ([]model.Channel, error) { + rows, err := s.DB.QueryContext(ctx, `SELECT id, name, type, config, status, created_at, updated_at FROM channel ORDER BY id`) + if err != nil { + return nil, fmt.Errorf("list channels: %w", err) + } + defer rows.Close() + + var channels []model.Channel + for rows.Next() { + var ch model.Channel + var configBytes []byte + if err := rows.Scan(&ch.ID, &ch.Name, &ch.Type, &configBytes, &ch.Status, &ch.CreatedAt, &ch.UpdatedAt); err != nil { + return nil, fmt.Errorf("scan channel: %w", err) + } + raw := json.RawMessage(configBytes) + ch.Config = &raw + channels = append(channels, ch) + } + return channels, rows.Err() +} + +func (s *Store) UpdateChannel(ctx context.Context, id int, ch *model.Channel) error { + configJSON, err := json.Marshal(ch.Config) + if err != nil { + return fmt.Errorf("marshal channel config: %w", err) + } + query := `UPDATE channel SET name=?, type=?, config=?, status=? WHERE id=?` + _, err = s.DB.ExecContext(ctx, query, ch.Name, ch.Type, configJSON, ch.Status, id) + if err != nil { + return fmt.Errorf("update channel %d: %w", id, err) + } + return nil +} + +func (s *Store) DeleteChannel(ctx context.Context, id int) error { + _, err := s.DB.ExecContext(ctx, `DELETE FROM channel WHERE id = ?`, id) + if err != nil { + return fmt.Errorf("delete channel %d: %w", id, err) + } + return nil +} +``` + +- [ ] **Step 3: Verify compilation** + +```bash +go build ./internal/store/... +``` +Expected: no errors. + +- [ ] **Step 4: Commit** + +```bash +git add -A && git commit -m "feat: template and channel store CRUD" +``` + +--- + +### Task 5: MySQL Store — Rule & MessageLog + +**Files:** +- Create: `internal/store/rule.go` +- Create: `internal/store/message_log.go` + +**Interfaces:** +- Produces: `Store.CreateRule`, `Store.GetRule`, `Store.GetRuleBySourceEvent`, `Store.ListRules`, `Store.UpdateRule`, `Store.DeleteRule`, `Store.AddRuleChannel`, `Store.RemoveRuleChannel`, `Store.GetRuleChannels`, `Store.UpdateRuleChannelEnabled` +- Produces: `Store.CreateMessageLog`, `Store.UpdateMessageLog`, `Store.ListMessageLogs`, `MessageLogFilter` + +- [ ] **Step 1: Write Rule store (includes RuleChannel ops)** + +File: `internal/store/rule.go` +```go +package store + +import ( + "context" + "encoding/json" + "fmt" + + "aiaa-notification-service/internal/model" +) + +func (s *Store) CreateRule(ctx context.Context, r *model.Rule, channelIDs []int) error { + tx, err := s.DB.BeginTxx(ctx, nil) + if err != nil { + return fmt.Errorf("begin tx: %w", err) + } + defer tx.Rollback() + + query := `INSERT INTO rule (source_id, event, template_id, conditions, enabled) VALUES (?, ?, ?, ?, ?)` + condsJSON, _ := marshalJSON(r.Conditions) + result, err := tx.ExecContext(ctx, query, r.SourceID, r.Event, r.TemplateID, condsJSON, r.Enabled) + if err != nil { + return fmt.Errorf("create rule: %w", err) + } + id, _ := result.LastInsertId() + r.ID = int(id) + + for _, chID := range channelIDs { + _, err := tx.ExecContext(ctx, `INSERT INTO rule_channel (rule_id, channel_id, enabled) VALUES (?, ?, 1)`, r.ID, chID) + if err != nil { + return fmt.Errorf("add rule_channel: %w", err) + } + } + return tx.Commit() +} + +func (s *Store) GetRule(ctx context.Context, id int) (*model.Rule, error) { + var r model.Rule + var condsBytes []byte + row := s.DB.QueryRowContext(ctx, `SELECT id, source_id, event, template_id, conditions, enabled, created_at, updated_at FROM rule WHERE id = ?`, id) + if err := row.Scan(&r.ID, &r.SourceID, &r.Event, &r.TemplateID, &condsBytes, &r.Enabled, &r.CreatedAt, &r.UpdatedAt); err != nil { + return nil, fmt.Errorf("get rule %d: %w", id, err) + } + if len(condsBytes) > 0 && string(condsBytes) != "null" { + raw := json.RawMessage(condsBytes) + r.Conditions = &raw + } + return &r, nil +} + +func (s *Store) GetRuleBySourceEvent(ctx context.Context, sourceID int, event string) (*model.Rule, error) { + var r model.Rule + var condsBytes []byte + query := `SELECT id, source_id, event, template_id, conditions, enabled, created_at, updated_at FROM rule WHERE source_id = ? AND event = ? AND enabled = 1` + row := s.DB.QueryRowContext(ctx, query, sourceID, event) + if err := row.Scan(&r.ID, &r.SourceID, &r.Event, &r.TemplateID, &condsBytes, &r.Enabled, &r.CreatedAt, &r.UpdatedAt); err != nil { + return nil, fmt.Errorf("get rule by source+event: %w", err) + } + if len(condsBytes) > 0 && string(condsBytes) != "null" { + raw := json.RawMessage(condsBytes) + r.Conditions = &raw + } + return &r, nil +} + +func (s *Store) ListRules(ctx context.Context) ([]model.Rule, error) { + rows, err := s.DB.QueryContext(ctx, `SELECT id, source_id, event, template_id, conditions, enabled, created_at, updated_at FROM rule ORDER BY id`) + if err != nil { + return nil, fmt.Errorf("list rules: %w", err) + } + defer rows.Close() + return scanRules(rows) +} + +func (s *Store) UpdateRule(ctx context.Context, id int, r *model.Rule, channelIDs []int) error { + tx, err := s.DB.BeginTxx(ctx, nil) + if err != nil { + return fmt.Errorf("begin tx: %w", err) + } + defer tx.Rollback() + + condsJSON, _ := marshalJSON(r.Conditions) + _, err = tx.ExecContext(ctx, `UPDATE rule SET source_id=?, event=?, template_id=?, conditions=?, enabled=? WHERE id=?`, + r.SourceID, r.Event, r.TemplateID, condsJSON, r.Enabled, id) + if err != nil { + return fmt.Errorf("update rule: %w", err) + } + + if channelIDs != nil { + _, _ = tx.ExecContext(ctx, `DELETE FROM rule_channel WHERE rule_id = ?`, id) + for _, chID := range channelIDs { + _, err := tx.ExecContext(ctx, `INSERT INTO rule_channel (rule_id, channel_id, enabled) VALUES (?, ?, 1)`, id, chID) + if err != nil { + return fmt.Errorf("add rule_channel: %w", err) + } + } + } + return tx.Commit() +} + +func (s *Store) DeleteRule(ctx context.Context, id int) error { + _, err := s.DB.ExecContext(ctx, `DELETE FROM rule WHERE id = ?`, id) + return err +} + +func (s *Store) SetRuleEnabled(ctx context.Context, id int, enabled bool) error { + v := 0 + if enabled { + v = 1 + } + _, err := s.DB.ExecContext(ctx, `UPDATE rule SET enabled = ? WHERE id = ?`, v, id) + return err +} + +func (s *Store) GetRuleChannels(ctx context.Context, ruleID int) ([]model.RuleChannel, error) { + var rcs []model.RuleChannel + err := s.DB.SelectContext(ctx, &rcs, `SELECT id, rule_id, channel_id, enabled FROM rule_channel WHERE rule_id = ? AND enabled = 1`, ruleID) + if err != nil { + return nil, fmt.Errorf("get rule channels: %w", err) + } + return rcs, nil +} + +func (s *Store) SetRuleChannelEnabled(ctx context.Context, ruleID, channelID int, enabled bool) error { + v := 0 + if enabled { + v = 1 + } + _, err := s.DB.ExecContext(ctx, `UPDATE rule_channel SET enabled = ? WHERE rule_id = ? AND channel_id = ?`, v, ruleID, channelID) + return err +} + +// helpers +func marshalJSON(v *json.RawMessage) (interface{}, error) { + if v == nil { + return nil, nil + } + return json.Marshal(v) +} + +func scanRules(rows *sql.Rows) ([]model.Rule, error) { + var rules []model.Rule + for rows.Next() { + var r model.Rule + var condsBytes []byte + if err := rows.Scan(&r.ID, &r.SourceID, &r.Event, &r.TemplateID, &condsBytes, &r.Enabled, &r.CreatedAt, &r.UpdatedAt); err != nil { + return nil, err + } + if len(condsBytes) > 0 && string(condsBytes) != "null" { + raw := json.RawMessage(condsBytes) + r.Conditions = &raw + } + rules = append(rules, r) + } + return rules, rows.Err() +} +``` + +Note: `scanRules` needs `database/sql` import. Add to top of file: +```go +import ( + "context" + "database/sql" + "encoding/json" + "fmt" + + "aiaa-notification-service/internal/model" +) +``` + +- [ ] **Step 2: Write MessageLog store** + +File: `internal/store/message_log.go` +```go +package store + +import ( + "context" + "fmt" + + "aiaa-notification-service/internal/model" +) + +type MessageLogFilter struct { + Source string `form:"source"` + Event string `form:"event"` + Status string `form:"status"` + Page int `form:"page"` + PageSize int `form:"page_size"` +} + +func (s *Store) CreateMessageLog(ctx context.Context, ml *model.MessageLog) error { + query := `INSERT INTO message_log (rule_id, channel_id, source, event, payload, content, status, retry_count) VALUES (?, ?, ?, ?, ?, ?, ?, ?)` + result, err := s.DB.ExecContext(ctx, query, ml.RuleID, ml.ChannelID, ml.Source, ml.Event, ml.Payload, ml.Content, ml.Status, ml.RetryCount) + if err != nil { + return fmt.Errorf("create message_log: %w", err) + } + id, _ := result.LastInsertId() + ml.ID = id + return nil +} + +func (s *Store) UpdateMessageLog(ctx context.Context, id int64, status string, response, errMsg *string) error { + query := `UPDATE message_log SET status=?, response=?, error_msg=? WHERE id=?` + _, err := s.DB.ExecContext(ctx, query, status, response, errMsg, id) + return err +} + +func (s *Store) ListMessageLogs(ctx context.Context, filter MessageLogFilter) ([]model.MessageLog, int, error) { + where := "WHERE 1=1" + args := []interface{}{} + if filter.Source != "" { + where += " AND source = ?" + args = append(args, filter.Source) + } + if filter.Event != "" { + where += " AND event = ?" + args = append(args, filter.Event) + } + if filter.Status != "" { + where += " AND status = ?" + args = append(args, filter.Status) + } + + var count int + countQuery := "SELECT COUNT(*) FROM message_log " + where + if err := s.DB.GetContext(ctx, &count, countQuery, args...); err != nil { + return nil, 0, err + } + + if filter.Page <= 0 { + filter.Page = 1 + } + if filter.PageSize <= 0 { + filter.PageSize = 20 + } + offset := (filter.Page - 1) * filter.PageSize + + var logs []model.MessageLog + query := "SELECT * FROM message_log " + where + " ORDER BY id DESC LIMIT ? OFFSET ?" + args = append(args, filter.PageSize, offset) + if err := s.DB.SelectContext(ctx, &logs, query, args...); err != nil { + return nil, 0, err + } + return logs, count, nil +} +``` + +- [ ] **Step 3: Verify compilation** + +```bash +go build ./internal/store/... +``` +Expected: no errors. Fix any import issues. + +- [ ] **Step 4: Commit** + +```bash +git add -A && git commit -m "feat: rule and message_log store CRUD" +``` + +--- + +### Task 6: Redis Cache Layer + +**Files:** +- Create: `internal/cache/redis.go` + +**Interfaces:** +- Consumes: `config.RedisConfig` +- Produces: `cache.NewCache(cfg config.RedisConfig) (*Cache, error)`, `Cache.GetRule`, `Cache.SetRule`, `Cache.InvalidateRule`, `Cache.GetChannels`, `Cache.SetChannels`, `Cache.InvalidateChannels`, `Cache.CheckRateLimit`, `Cache.InvalidateBySource`, `Cache.InvalidateByTemplate` + +- [ ] **Step 1: Write cache layer** + +File: `internal/cache/redis.go` +```go +package cache + +import ( + "context" + "encoding/json" + "fmt" + "strconv" + "time" + + "aiaa-notification-service/internal/model" + + "github.com/redis/go-redis/v9" +) + +type Cache struct { + rdb *redis.Client +} + +// CachedRule holds the minimal info needed after rule matching. +type CachedRule struct { + RuleID int `json:"rule_id"` + TemplateID int `json:"template_id"` + Content string `json:"content"` + Conditions string `json:"conditions"` // JSON string, empty if null +} + +type CachedChannel struct { + ID int `json:"id"` + Type string `json:"type"` + Config *json.RawMessage `json:"config"` +} + +func NewCache(cfg model.config.RedisConfig) (*Cache, error) { + rdb := redis.NewClient(&redis.Options{ + Addr: cfg.Addr(), + Password: cfg.Password, + DB: cfg.DB, + }) + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + defer cancel() + if err := rdb.Ping(ctx).Err(); err != nil { + return nil, fmt.Errorf("redis ping: %w", err) + } + return &Cache{rdb: rdb}, nil +} + +// --- Rule cache --- + +func ruleKey(sourceID int, event string) string { + return fmt.Sprintf("notify:rule:%d:%s", sourceID, event) +} + +func (c *Cache) GetRule(ctx context.Context, sourceID int, event string) (*CachedRule, error) { + data, err := c.rdb.Get(ctx, ruleKey(sourceID, event)).Bytes() + if err != nil { + return nil, err + } + var cr CachedRule + if err := json.Unmarshal(data, &cr); err != nil { + return nil, err + } + return &cr, nil +} + +func (c *Cache) SetRule(ctx context.Context, sourceID int, event string, cr *CachedRule) error { + data, err := json.Marshal(cr) + if err != nil { + return err + } + return c.rdb.Set(ctx, ruleKey(sourceID, event), data, 5*time.Minute).Err() +} + +func (c *Cache) InvalidateRule(ctx context.Context, sourceID int, event string) error { + return c.rdb.Del(ctx, ruleKey(sourceID, event)).Err() +} + +// --- Channel cache --- + +func channelsKey(ruleID int) string { + return fmt.Sprintf("notify:channels:%d", ruleID) +} + +func (c *Cache) GetChannels(ctx context.Context, ruleID int) ([]CachedChannel, error) { + data, err := c.rdb.Get(ctx, channelsKey(ruleID)).Bytes() + if err != nil { + return nil, err + } + var channels []CachedChannel + if err := json.Unmarshal(data, &channels); err != nil { + return nil, err + } + return channels, nil +} + +func (c *Cache) SetChannels(ctx context.Context, ruleID int, channels []CachedChannel) error { + data, err := json.Marshal(channels) + if err != nil { + return err + } + return c.rdb.Set(ctx, channelsKey(ruleID), data, 5*time.Minute).Err() +} + +func (c *Cache) InvalidateChannels(ctx context.Context, ruleID int) error { + return c.rdb.Del(ctx, channelsKey(ruleID)).Err() +} + +// --- Bulk invalidation --- + +func (c *Cache) InvalidateBySource(ctx context.Context, sourceID int) error { + pattern := fmt.Sprintf("notify:rule:%d:*", sourceID) + iter := c.rdb.Scan(ctx, 0, pattern, 0).Iterator() + for iter.Next(ctx) { + c.rdb.Del(ctx, iter.Val()) + } + return iter.Err() +} + +func (c *Cache) InvalidateByTemplate(ctx context.Context, templateID int) error { + // Template changes mean all cached rule info is stale. Simplest: scan all rule keys. + // In production, you'd maintain a template→rule reverse index. For v1, scan is acceptable. + iter := c.rdb.Scan(ctx, 0, "notify:rule:*", 0).Iterator() + for iter.Next(ctx) { + c.rdb.Del(ctx, iter.Val()) + } + return iter.Err() +} + +// --- Rate limiter --- + +func (c *Cache) CheckRateLimit(ctx context.Context, sourceID int, limitPerSec int) (bool, int, error) { + key := fmt.Sprintf("ratelimit:%d", sourceID) + now := time.Now().Unix() + pipe := c.rdb.Pipeline() + pipe.ZRemRangeByScore(ctx, key, "0", strconv.FormatInt(now-1, 10)) + countCmd := pipe.ZCard(ctx, key) + pipe.Exec(ctx) + + count := countCmd.Val() + if count >= int64(limitPerSec) { + return false, 1, nil // rate limited, retry after 1s + } + + c.rdb.ZAdd(ctx, key, redis.Z{Score: float64(now), Member: strconv.FormatInt(now*1000, 10)}) + c.rdb.Expire(ctx, key, 2*time.Second) + return true, 0, nil +} + +func (c *Cache) Close() error { + return c.rdb.Close() +} +``` + +Note: the import references `model.config.RedisConfig` — need to fix this. The cache package imports `config`, not model. Update the function signature: + +```go +import ( + "aiaa-notification-service/internal/config" +) + +func NewCache(cfg config.RedisConfig) (*Cache, error) { +``` + +- [ ] **Step 2: Install go-redis** + +```bash +go get github.com/redis/go-redis/v9 +go mod tidy +``` + +- [ ] **Step 3: Verify compilation** + +```bash +go build ./internal/cache/... +``` +Expected: no errors. + +- [ ] **Step 4: Commit** + +```bash +git add -A && git commit -m "feat: redis cache layer with rule/channel caching and rate limiter" +``` + +--- + +### Task 7: Message Parser (JSON / Regex / Text) + +**Files:** +- Create: `internal/parser/parser.go` +- Create: `internal/parser/json.go` +- Create: `internal/parser/regex.go` +- Create: `internal/parser/text.go` + +**Interfaces:** +- Produces: `parser.ParsedMessage`, `parser.Parser` interface, `parser.NewParser(parseMode string, parsePattern string) (Parser, error)` + +- [ ] **Step 1: Write parser interface and ParsedMessage** + +File: `internal/parser/parser.go` +```go +package parser + +import ( + "fmt" +) + +// ParsedMessage is the result of parsing an incoming notification body. +type ParsedMessage struct { + Event string + Data map[string]interface{} // field name → value, passed to template + Body string // raw body (for text mode) +} + +// Parser converts a raw HTTP body into a ParsedMessage. +type Parser interface { + Parse(body []byte) (*ParsedMessage, error) +} + +func NewParser(parseMode string, parsePattern string) (Parser, error) { + switch parseMode { + case "json": + return &JSONParser{}, nil + case "regex": + if parsePattern == "" { + return nil, fmt.Errorf("regex parse mode requires parse_pattern") + } + return NewRegexParser(parsePattern) + case "text": + return &TextParser{}, nil + default: + return nil, fmt.Errorf("unknown parse_mode: %s", parseMode) + } +} +``` + +- [ ] **Step 2: Write JSON parser** + +File: `internal/parser/json.go` +```go +package parser + +import ( + "encoding/json" + "fmt" +) + +type JSONParser struct{} + +func (p *JSONParser) Parse(body []byte) (*ParsedMessage, error) { + var raw map[string]interface{} + if err := json.Unmarshal(body, &raw); err != nil { + return nil, fmt.Errorf("json parse: %w", err) + } + + event, ok := raw["event"].(string) + if !ok { + return nil, fmt.Errorf("json: missing 'event' field") + } + + data, _ := raw["data"].(map[string]interface{}) + if data == nil { + // If no "data" key, use the entire JSON as data (minus event) + data = make(map[string]interface{}) + for k, v := range raw { + if k != "event" { + data[k] = v + } + } + } + + return &ParsedMessage{ + Event: event, + Data: data, + Body: string(body), + }, nil +} +``` + +- [ ] **Step 3: Write Regex parser** + +File: `internal/parser/regex.go` +```go +package parser + +import ( + "fmt" + "regexp" +) + +type RegexParser struct { + re *regexp.Regexp +} + +func NewRegexParser(pattern string) (*RegexParser, error) { + re, err := regexp.Compile(pattern) + if err != nil { + return nil, fmt.Errorf("compile regex pattern: %w", err) + } + return &RegexParser{re: re}, nil +} + +func (p *RegexParser) Parse(body []byte) (*ParsedMessage, error) { + matches := p.re.FindStringSubmatch(string(body)) + if matches == nil { + return nil, fmt.Errorf("regex: body does not match pattern") + } + + names := p.re.SubexpNames() + data := make(map[string]interface{}) + var event string + + for i, name := range names { + if i == 0 || name == "" { + continue + } + if name == "event" { + event = matches[i] + } else { + data[name] = matches[i] + } + } + + return &ParsedMessage{ + Event: event, + Data: data, + Body: string(body), + }, nil +} +``` + +- [ ] **Step 4: Write Text parser** + +File: `internal/parser/text.go` +```go +package parser + +type TextParser struct{} + +func (p *TextParser) Parse(body []byte) (*ParsedMessage, error) { + return &ParsedMessage{ + Event: "default", // text mode has no explicit event + Data: map[string]interface{}{"Body": string(body)}, + Body: string(body), + }, nil +} +``` + +- [ ] **Step 5: Write parser tests** + +File: `internal/parser/parser_test.go` +```go +package parser + +import ( + "testing" +) + +func TestJSONParser(t *testing.T) { + p := &JSONParser{} + body := []byte(`{"event":"trade.open","data":{"symbol":"BTC","price":65000}}`) + msg, err := p.Parse(body) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if msg.Event != "trade.open" { + t.Errorf("expected event 'trade.open', got '%s'", msg.Event) + } + if msg.Data["symbol"] != "BTC" { + t.Errorf("expected symbol 'BTC', got '%v'", msg.Data["symbol"]) + } +} + +func TestJSONParser_NoDataKey(t *testing.T) { + p := &JSONParser{} + body := []byte(`{"event":"alert","symbol":"ETH"}`) + msg, err := p.Parse(body) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if msg.Data["symbol"] != "ETH" { + t.Errorf("expected symbol 'ETH' in data, got '%v'", msg.Data["symbol"]) + } +} + +func TestRegexParser(t *testing.T) { + p, err := NewRegexParser(`(?P\w+) (?P.+)`) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + body := []byte("trade.open BTC long at 65000") + msg, err := p.Parse(body) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if msg.Event != "trade.open" { + t.Errorf("expected event 'trade.open', got '%s'", msg.Event) + } + if msg.Data["message"] != "BTC long at 65000" { + t.Errorf("expected message, got '%v'", msg.Data["message"]) + } +} + +func TestTextParser(t *testing.T) { + p := &TextParser{} + body := []byte("raw notification text here") + msg, err := p.Parse(body) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if msg.Data["Body"] != "raw notification text here" { + t.Errorf("expected Body in data, got '%v'", msg.Data["Body"]) + } +} +``` + +- [ ] **Step 6: Run tests** + +```bash +go test ./internal/parser/... -v +``` +Expected: all tests PASS. + +- [ ] **Step 7: Commit** + +```bash +git add -A && git commit -m "feat: message parser with json/regex/text modes" +``` + +--- + +### Task 8: Condition Evaluator + +**Files:** +- Create: `internal/condition/evaluator.go` + +**Interfaces:** +- Produces: `condition.Evaluate(conditions []model.Condition, data map[string]interface{}) bool` + +- [ ] **Step 1: Write condition evaluator** + +File: `internal/condition/evaluator.go` +```go +package condition + +import ( + "fmt" + "strconv" + "strings" + + "aiaa-notification-service/internal/model" +) + +// Evaluate checks all conditions against data. Returns true if ALL conditions pass. +// An empty conditions slice always returns true. +func Evaluate(conditions []model.Condition, data map[string]interface{}) bool { + if len(conditions) == 0 { + return true + } + for _, c := range conditions { + if !evaluateOne(c, data) { + return false + } + } + return true +} + +func evaluateOne(c model.Condition, data map[string]interface{}) bool { + fieldVal, fieldExists := data[c.Field] + + switch c.Op { + case "exists": + return fieldExists + case "not_exists": + return !fieldExists + case "eq": + if !fieldExists { + return false + } + return fmt.Sprintf("%v", fieldVal) == c.Value + case "ne": + if !fieldExists { + return false + } + return fmt.Sprintf("%v", fieldVal) != c.Value + case "contains": + if !fieldExists { + return false + } + return strings.Contains(fmt.Sprintf("%v", fieldVal), c.Value) + case "gt", "gte", "lt", "lte": + if !fieldExists { + return false + } + return compareNumeric(fieldVal, c.Value, c.Op) + default: + return false + } +} + +func compareNumeric(fieldVal interface{}, value string, op string) bool { + fv, err := toFloat64(fieldVal) + if err != nil { + return false + } + cv, err := strconv.ParseFloat(value, 64) + if err != nil { + return false + } + switch op { + case "gt": + return fv > cv + case "gte": + return fv >= cv + case "lt": + return fv < cv + case "lte": + return fv <= cv + } + return false +} + +func toFloat64(v interface{}) (float64, error) { + switch val := v.(type) { + case float64: + return val, nil + case float32: + return float64(val), nil + case int: + return float64(val), nil + case int64: + return float64(val), nil + case string: + return strconv.ParseFloat(val, 64) + case json.Number: + return val.Float64() + default: + return 0, fmt.Errorf("cannot convert %T to float64", v) + } +} +``` + +Note: needs `encoding/json` import for `json.Number`. + +- [ ] **Step 2: Write tests** + +File: `internal/condition/evaluator_test.go` +```go +package condition + +import ( + "testing" + + "aiaa-notification-service/internal/model" +) + +func TestEvaluate_Empty(t *testing.T) { + if !Evaluate(nil, nil) { + t.Error("empty conditions should pass") + } +} + +func TestEvaluate_Exists(t *testing.T) { + conds := []model.Condition{{Field: "symbol", Op: "exists"}} + data := map[string]interface{}{"symbol": "BTC"} + if !Evaluate(conds, data) { + t.Error("symbol exists, should pass") + } +} + +func TestEvaluate_NotExists(t *testing.T) { + conds := []model.Condition{{Field: "symbol", Op: "not_exists"}} + data := map[string]interface{}{"price": 100} + if !Evaluate(conds, data) { + t.Error("symbol not exists, should pass") + } +} + +func TestEvaluate_Gt(t *testing.T) { + conds := []model.Condition{{Field: "price", Op: "gt", Value: "100"}} + data := map[string]interface{}{"price": float64(200)} + if !Evaluate(conds, data) { + t.Error("200 > 100, should pass") + } +} + +func TestEvaluate_Fail(t *testing.T) { + conds := []model.Condition{ + {Field: "symbol", Op: "exists"}, + {Field: "price", Op: "lt", Value: "100"}, + } + data := map[string]interface{}{"symbol": "BTC", "price": float64(200)} + if Evaluate(conds, data) { + t.Error("200 < 100 is false, should fail") + } +} + +func TestEvaluate_Contains(t *testing.T) { + conds := []model.Condition{{Field: "msg", Op: "contains", Value: "error"}} + data := map[string]interface{}{"msg": "connection error occurred"} + if !Evaluate(conds, data) { + t.Error("msg contains 'error', should pass") + } +} +``` + +- [ ] **Step 3: Run tests** + +```bash +go test ./internal/condition/... -v +``` +Expected: all tests PASS. + +- [ ] **Step 4: Commit** + +```bash +git add -A && git commit -m "feat: condition evaluator for rule filtering" +``` + +--- + +### Task 9: Channel Adapters — Interface + HTTP-based (DingTalk, WeCom, Bark) + +**Files:** +- Create: `internal/adapter/adapter.go` +- Create: `internal/adapter/dingtalk.go` +- Create: `internal/adapter/wecom.go` +- Create: `internal/adapter/bark.go` + +**Interfaces:** +- Produces: `adapter.ChannelSender` interface, `adapter.NewSender(channelType string, smtpCfg *config.SMTPConfig) (ChannelSender, error)` + +- [ ] **Step 1: Write ChannelSender interface + factory** + +File: `internal/adapter/adapter.go` +```go +package adapter + +import ( + "encoding/json" + "fmt" + + "aiaa-notification-service/internal/config" +) + +// ChannelSender sends a rendered message to a specific channel. +type ChannelSender interface { + Type() string + Send(title, content string, config json.RawMessage) error +} + +func NewSender(channelType string, smtpCfg *config.SMTPConfig) (ChannelSender, error) { + switch channelType { + case "dingtalk": + return &DingTalkSender{}, nil + case "wecom": + return &WeComSender{}, nil + case "email": + if smtpCfg == nil { + return nil, fmt.Errorf("SMTP config required for email sender") + } + return NewEmailSender(*smtpCfg), nil + case "bark": + return &BarkSender{}, nil + default: + return nil, fmt.Errorf("unknown channel type: %s", channelType) + } +} +``` + +- [ ] **Step 2: Write DingTalk sender** + +File: `internal/adapter/dingtalk.go` +```go +package adapter + +import ( + "bytes" + "crypto/hmac" + "crypto/sha256" + "encoding/base64" + "encoding/json" + "fmt" + "net/http" + "net/url" + "time" +) + +type dingtalkConfig struct { + WebhookURL string `json:"webhook_url"` + Secret string `json:"secret,omitempty"` +} + +type dingtalkMessage struct { + MsgType string `json:"msgtype"` + Markdown *dingtalkMD `json:"markdown"` +} + +type dingtalkMD struct { + Title string `json:"title"` + Text string `json:"text"` +} + +type DingTalkSender struct{} + +func (s *DingTalkSender) Type() string { return "dingtalk" } + +func (s *DingTalkSender) Send(title, content string, config json.RawMessage) error { + var cfg dingtalkConfig + if err := json.Unmarshal(config, &cfg); err != nil { + return fmt.Errorf("parse dingtalk config: %w", err) + } + + reqURL := cfg.WebhookURL + if cfg.Secret != "" { + timestamp := time.Now().UnixMilli() + sign := dingtalkSign(timestamp, cfg.Secret) + reqURL = fmt.Sprintf("%s×tamp=%d&sign=%s", cfg.WebhookURL, timestamp, sign) + } + + msg := dingtalkMessage{ + MsgType: "markdown", + Markdown: &dingtalkMD{ + Title: title, + Text: content, + }, + } + + body, _ := json.Marshal(msg) + resp, err := http.Post(reqURL, "application/json", bytes.NewReader(body)) + if err != nil { + return fmt.Errorf("dingtalk send: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode >= 400 { + return fmt.Errorf("dingtalk returned status %d", resp.StatusCode) + } + return nil +} + +func dingtalkSign(timestamp int64, secret string) string { + mac := hmac.New(sha256.New, []byte(secret)) + fmt.Fprintf(mac, "%d\n%s", timestamp, secret) + signData := mac.Sum(nil) + return url.QueryEscape(base64.StdEncoding.EncodeToString(signData)) +} +``` + +- [ ] **Step 3: Write WeCom (企业微信) sender** + +File: `internal/adapter/wecom.go` +```go +package adapter + +import ( + "bytes" + "encoding/json" + "fmt" + "net/http" +) + +type wecomConfig struct { + WebhookURL string `json:"webhook_url"` +} + +type wecomMessage struct { + MsgType string `json:"msgtype"` + Markdown *wecomMD `json:"markdown"` +} + +type wecomMD struct { + Content string `json:"content"` +} + +type WeComSender struct{} + +func (s *WeComSender) Type() string { return "wecom" } + +func (s *WeComSender) Send(title, content string, config json.RawMessage) error { + var cfg wecomConfig + if err := json.Unmarshal(config, &cfg); err != nil { + return fmt.Errorf("parse wecom config: %w", err) + } + + msg := wecomMessage{ + MsgType: "markdown", + Markdown: &wecomMD{ + Content: fmt.Sprintf("## %s\n%s", title, content), + }, + } + + body, _ := json.Marshal(msg) + resp, err := http.Post(cfg.WebhookURL, "application/json", bytes.NewReader(body)) + if err != nil { + return fmt.Errorf("wecom send: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode >= 400 { + return fmt.Errorf("wecom returned status %d", resp.StatusCode) + } + return nil +} +``` + +- [ ] **Step 4: Write Bark sender** + +File: `internal/adapter/bark.go` +```go +package adapter + +import ( + "bytes" + "encoding/json" + "fmt" + "net/http" +) + +type barkConfig struct { + URL string `json:"url"` +} + +type barkPayload struct { + Title string `json:"title"` + Body string `json:"body"` +} + +type BarkSender struct{} + +func (s *BarkSender) Type() string { return "bark" } + +func (s *BarkSender) Send(title, content string, config json.RawMessage) error { + var cfg barkConfig + if err := json.Unmarshal(config, &cfg); err != nil { + return fmt.Errorf("parse bark config: %w", err) + } + + payload := barkPayload{Title: title, Body: content} + body, _ := json.Marshal(payload) + resp, err := http.Post(cfg.URL, "application/json", bytes.NewReader(body)) + if err != nil { + return fmt.Errorf("bark send: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode >= 400 { + return fmt.Errorf("bark returned status %d", resp.StatusCode) + } + return nil +} +``` + +- [ ] **Step 5: Verify compilation** + +```bash +go build ./internal/adapter/... +``` +Expected: no errors. + +- [ ] **Step 6: Commit** + +```bash +git add -A && git commit -m "feat: channel adapters — dingtalk, wecom, bark" +``` + +--- + +### Task 10: Email Adapter (SMTP) + +**Files:** +- Create: `internal/adapter/email.go` + +**Interfaces:** +- Produces: `EmailSender` implementing `ChannelSender` + +- [ ] **Step 1: Write Email sender** + +File: `internal/adapter/email.go` +```go +package adapter + +import ( + "crypto/tls" + "encoding/json" + "fmt" + "net" + "net/smtp" + "strings" + + "aiaa-notification-service/internal/config" +) + +type emailConfig struct { + To []string `json:"to"` +} + +type EmailSender struct { + cfg config.SMTPConfig +} + +func NewEmailSender(cfg config.SMTPConfig) *EmailSender { + return &EmailSender{cfg: cfg} +} + +func (s *EmailSender) Type() string { return "email" } + +func (s *EmailSender) Send(title, content string, config json.RawMessage) error { + var ecfg emailConfig + if err := json.Unmarshal(config, &ecfg); err != nil { + return fmt.Errorf("parse email config: %w", err) + } + if len(ecfg.To) == 0 { + return fmt.Errorf("email: 'to' list is empty") + } + + msg := buildEmail(s.cfg.From, ecfg.To, title, content) + addr := fmt.Sprintf("%s:%d", s.cfg.Host, s.cfg.Port) + + if s.cfg.Port == 465 { + return s.sendTLS(addr, ecfg.To, msg) + } + return s.sendSTARTTLS(addr, ecfg.To, msg) +} + +func (s *EmailSender) sendTLS(addr string, to []string, msg []byte) error { + tlsCfg := &tls.Config{ServerName: s.cfg.Host} + conn, err := tls.Dial("tcp", addr, tlsCfg) + if err != nil { + return fmt.Errorf("tls dial: %w", err) + } + defer conn.Close() + + client, err := smtp.NewClient(conn, s.cfg.Host) + if err != nil { + return fmt.Errorf("smtp client: %w", err) + } + defer client.Quit() + + return s.authAndSend(client, to, msg) +} + +func (s *EmailSender) sendSTARTTLS(addr string, to []string, msg []byte) error { + conn, err := net.Dial("tcp", addr) + if err != nil { + return fmt.Errorf("dial: %w", err) + } + defer conn.Close() + + client, err := smtp.NewClient(conn, s.cfg.Host) + if err != nil { + return fmt.Errorf("smtp client: %w", err) + } + defer client.Quit() + + if ok, _ := client.Extension("STARTTLS"); ok { + tlsCfg := &tls.Config{ServerName: s.cfg.Host} + if err := client.StartTLS(tlsCfg); err != nil { + return fmt.Errorf("starttls: %w", err) + } + } + + return s.authAndSend(client, to, msg) +} + +func (s *EmailSender) authAndSend(client *smtp.Client, to []string, msg []byte) error { + if s.cfg.User != "" { + auth := smtp.PlainAuth("", s.cfg.User, s.cfg.Password, s.cfg.Host) + if err := client.Auth(auth); err != nil { + return fmt.Errorf("auth: %w", err) + } + } + + if err := client.Mail(s.cfg.From); err != nil { + return fmt.Errorf("mail: %w", err) + } + for _, recipient := range to { + if err := client.Rcpt(recipient); err != nil { + return fmt.Errorf("rcpt %s: %w", recipient, err) + } + } + w, err := client.Data() + if err != nil { + return fmt.Errorf("data: %w", err) + } + _, err = w.Write(msg) + if err != nil { + return fmt.Errorf("write: %w", err) + } + return w.Close() +} + +func buildEmail(from string, to []string, subject, body string) []byte { + var sb strings.Builder + sb.WriteString(fmt.Sprintf("From: %s\r\n", from)) + sb.WriteString(fmt.Sprintf("To: %s\r\n", strings.Join(to, ", "))) + sb.WriteString(fmt.Sprintf("Subject: %s\r\n", subject)) + sb.WriteString("MIME-Version: 1.0\r\n") + sb.WriteString("Content-Type: text/plain; charset=utf-8\r\n") + sb.WriteString("\r\n") + sb.WriteString(body) + return []byte(sb.String()) +} +``` + +Note: fix `s.cfg` references — they should be `s.cfg`. + +- [ ] **Step 2: Verify compilation** + +```bash +go build ./internal/adapter/... +``` +Expected: no errors. + +- [ ] **Step 3: Commit** + +```bash +git add -A && git commit -m "feat: email adapter with SMTP TLS support" +``` + +--- + +### Task 11: Engine — Renderer, Matcher, Router + +**Files:** +- Create: `internal/engine/renderer.go` +- Create: `internal/engine/matcher.go` +- Create: `internal/engine/router.go` + +**Interfaces:** +- Produces: `engine.NewRenderer() *Renderer`, `Renderer.Render(templateContent string, data map[string]interface{}) (string, error)` +- Produces: `engine.NewMatcher(store *store.Store, cache *cache.Cache) *Matcher`, `Matcher.Match(ctx, sourceID, event) (*model.Rule, error)` +- Produces: `engine.NewRouter(store *store.Store, cache *cache.Cache, senderFactory func(string) (adapter.ChannelSender, error)) *Router`, `Router.Route(ctx, rule *model.Rule, content string) []SendResult` + +- [ ] **Step 1: Write Template Renderer** + +File: `internal/engine/renderer.go` +```go +package engine + +import ( + "bytes" + "fmt" + "text/template" +) + +type Renderer struct{} + +func NewRenderer() *Renderer { + return &Renderer{} +} + +func (r *Renderer) Render(tmplContent string, data map[string]interface{}) (string, error) { + tmpl, err := template.New("notify").Parse(tmplContent) + if err != nil { + return "", fmt.Errorf("parse template: %w", err) + } + var buf bytes.Buffer + if err := tmpl.Execute(&buf, data); err != nil { + return "", fmt.Errorf("execute template: %w", err) + } + return buf.String(), nil +} +``` + +- [ ] **Step 2: Write Rule Matcher** + +File: `internal/engine/matcher.go` +```go +package engine + +import ( + "context" + "encoding/json" + "fmt" + + "aiaa-notification-service/internal/cache" + "aiaa-notification-service/internal/model" + "aiaa-notification-service/internal/store" +) + +type Matcher struct { + store *store.Store + cache *cache.Cache +} + +func NewMatcher(s *store.Store, c *cache.Cache) *Matcher { + return &Matcher{store: s, cache: c} +} + +func (m *Matcher) Match(ctx context.Context, sourceID int, event string) (*model.Rule, error) { + // Try cache first + if m.cache != nil { + cr, err := m.cache.GetRule(ctx, sourceID, event) + if err == nil { + rule := &model.Rule{ID: cr.RuleID, TemplateID: cr.TemplateID, SourceID: sourceID, Event: event} + if cr.Conditions != "" && cr.Conditions != "null" { + raw := json.RawMessage(cr.Conditions) + rule.Conditions = &raw + } + return rule, nil + } + } + + // Fall back to DB + rule, err := m.store.GetRuleBySourceEvent(ctx, sourceID, event) + if err != nil { + return nil, fmt.Errorf("match rule: %w", err) + } + + // Warm cache + if m.cache != nil { + tmpl, err := m.store.GetTemplate(ctx, rule.TemplateID) + if err != nil { + return rule, nil // rule found but template fetch failed — still return rule + } + cr := &cache.CachedRule{ + RuleID: rule.ID, + TemplateID: rule.TemplateID, + Content: tmpl.Content, + } + if rule.Conditions != nil { + cr.Conditions = string(*rule.Conditions) + } + _ = m.cache.SetRule(ctx, sourceID, event, cr) + } + + return rule, nil +} +``` + +- [ ] **Step 3: Write Router** + +File: `internal/engine/router.go` +```go +package engine + +import ( + "context" + "fmt" + "log/slog" + + "aiaa-notification-service/internal/adapter" + "aiaa-notification-service/internal/cache" + "aiaa-notification-service/internal/model" + "aiaa-notification-service/internal/store" +) + +type SendResult struct { + ChannelType string `json:"channel_type"` + ChannelID int `json:"channel_id"` + Success bool `json:"success"` + Error string `json:"error,omitempty"` +} + +type Router struct { + store *store.Store + cache *cache.Cache + senderFactory func(channelType string) (adapter.ChannelSender, error) +} + +func NewRouter(s *store.Store, c *cache.Cache, sf func(channelType string) (adapter.ChannelSender, error)) *Router { + return &Router{store: s, cache: c, senderFactory: sf} +} + +// Route sends content to all enabled channels for the rule. Returns immediately, sends async. +func (r *Router) Route(ctx context.Context, rule *model.Rule, title, content string) []string { + channels, err := r.getChannels(ctx, rule.ID) + if err != nil { + slog.Error("get channels for rule", "rule_id", rule.ID, "error", err) + return nil + } + + channelNames := make([]string, 0, len(channels)) + for _, ch := range channels { + channelNames = append(channelNames, fmt.Sprintf("%s:%d", ch.Type, ch.ID)) + sender, err := r.senderFactory(ch.Type) + if err != nil { + slog.Error("create sender", "type", ch.Type, "error", err) + continue + } + + go func(ch cache.CachedChannel, s adapter.ChannelSender) { + var cfg json.RawMessage + if ch.Config != nil { + cfg = *ch.Config + } + if err := s.Send(title, content, cfg); err != nil { + slog.Error("send failed", "channel_type", ch.Type, "channel_id", ch.ID, "error", err) + } else { + slog.Info("sent", "channel_type", ch.Type, "channel_id", ch.ID) + } + }(ch, sender) + } + + return channelNames +} + +func (r *Router) getChannels(ctx context.Context, ruleID int) ([]cache.CachedChannel, error) { + if r.cache != nil { + chs, err := r.cache.GetChannels(ctx, ruleID) + if err == nil { + return chs, nil + } + } + + rcs, err := r.store.GetRuleChannels(ctx, ruleID) + if err != nil { + return nil, err + } + + var cached []cache.CachedChannel + for _, rc := range rcs { + ch, err := r.store.GetChannel(ctx, rc.ChannelID) + if err != nil { + continue + } + cc := cache.CachedChannel{ID: ch.ID, Type: ch.Type, Config: ch.Config} + cached = append(cached, cc) + } + + if r.cache != nil { + _ = r.cache.SetChannels(ctx, ruleID, cached) + } + return cached, nil +} +``` + +Note: need `encoding/json` import in router.go. + +- [ ] **Step 4: Write tests for renderer** + +File: `internal/engine/renderer_test.go` +```go +package engine + +import ( + "strings" + "testing" +) + +func TestRenderer(t *testing.T) { + r := NewRenderer() + tmpl := "🚀 {{.symbol}} 开仓通知\n价格: {{.price}}" + data := map[string]interface{}{"symbol": "BTC", "price": 65000} + result, err := r.Render(tmpl, data) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if !strings.Contains(result, "BTC") || !strings.Contains(result, "65000") { + t.Errorf("unexpected output: %s", result) + } +} + +func TestRenderer_Error(t *testing.T) { + r := NewRenderer() + _, err := r.Render("{{.nonexistent}}", map[string]interface{}{}) + if err == nil { + t.Error("expected error for missing field, got nil") + } +} +``` + +- [ ] **Step 5: Run tests** + +```bash +go test ./internal/engine/... -v +``` +Expected: renderer tests PASS. + +- [ ] **Step 6: Commit** + +```bash +git add -A && git commit -m "feat: engine — template renderer, rule matcher, channel router" +``` + +--- + +### Task 12: HTTP Middleware (Auth + Rate Limit) + +**Files:** +- Create: `internal/handler/middleware.go` + +**Interfaces:** +- Produces: `handler.AdminAuth(adminKey string) gin.HandlerFunc`, `handler.SourceAuth(store *store.Store) gin.HandlerFunc`, `handler.RateLimit(cache *cache.Cache, defaultLimit int) gin.HandlerFunc` + +- [ ] **Step 1: Write middleware** + +File: `internal/handler/middleware.go` +```go +package handler + +import ( + "net/http" + "strings" + + "aiaa-notification-service/internal/cache" + "aiaa-notification-service/internal/store" + + "github.com/gin-gonic/gin" +) + +// AdminAuth checks the admin key for management API endpoints. +func AdminAuth(adminKey string) gin.HandlerFunc { + return func(c *gin.Context) { + key := extractBearer(c) + if key == "" || key != adminKey { + c.AbortWithStatusJSON(http.StatusUnauthorized, gin.H{"error": "unauthorized"}) + return + } + c.Next() + } +} + +// SourceAuth identifies the source by its API key and sets it in context. +func SourceAuth(s *store.Store) gin.HandlerFunc { + return func(c *gin.Context) { + key := extractBearer(c) + if key == "" { + c.AbortWithStatusJSON(http.StatusUnauthorized, gin.H{"error": "missing api key"}) + return + } + source, err := s.GetSourceByAPIKey(c.Request.Context(), key) + if err != nil { + c.AbortWithStatusJSON(http.StatusUnauthorized, gin.H{"error": "invalid api key"}) + return + } + c.Set("source", source) + c.Next() + } +} + +// RateLimit applies per-source rate limiting. +func RateLimit(cache *cache.Cache, defaultLimit int) gin.HandlerFunc { + return func(c *gin.Context) { + source, exists := c.Get("source") + if !exists { + c.Next() + return + } + src := source.(*model.Source) + allowed, retryAfter, err := cache.CheckRateLimit(c.Request.Context(), src.ID, defaultLimit) + if err != nil { + // Redis error — allow pass through + c.Next() + return + } + if !allowed { + c.AbortWithStatusJSON(http.StatusTooManyRequests, gin.H{ + "error": "rate_limit_exceeded", + "message": "too many requests", + "retry_after": retryAfter, + }) + return + } + c.Next() + } +} + +func extractBearer(c *gin.Context) string { + auth := c.GetHeader("Authorization") + if auth == "" { + return "" + } + parts := strings.SplitN(auth, " ", 2) + if len(parts) != 2 || !strings.EqualFold(parts[0], "bearer") { + return "" + } + return parts[1] +} +``` + +Note: `rateLimit` needs the `model` import for `source.(*model.Source)`. Add: +```go +import ( + "aiaa-notification-service/internal/model" +) +``` + +- [ ] **Step 2: Verify compilation** + +```bash +go get github.com/gin-gonic/gin +go build ./internal/handler/... +``` +Expected: no errors. + +- [ ] **Step 3: Commit** + +```bash +git add -A && git commit -m "feat: auth and rate limit middleware" +``` + +--- + +### Task 13: Notify Handler (POST /api/v1/notify) + +**Files:** +- Create: `internal/handler/notify.go` + +**Interfaces:** +- Consumes: `*store.Store`, `*cache.Cache`, `*engine.Matcher`, `*engine.Renderer`, `*engine.Router`, `*condition.Evaluator` +- Produces: `POST /api/v1/notify` endpoint + +- [ ] **Step 1: Write notify handler** + +File: `internal/handler/notify.go` +```go +package handler + +import ( + "encoding/json" + "io" + "log/slog" + "net/http" + + "aiaa-notification-service/internal/condition" + "aiaa-notification-service/internal/engine" + "aiaa-notification-service/internal/model" + "aiaa-notification-service/internal/parser" + "aiaa-notification-service/internal/store" + + "github.com/gin-gonic/gin" +) + +type NotifyHandler struct { + store *store.Store + cache *cache.Cache // used by engine + matcher *engine.Matcher + renderer *engine.Renderer + router *engine.Router +} + +func NewNotifyHandler(s *store.Store, c *cache.Cache, m *engine.Matcher, r *engine.Renderer, rt *engine.Router) *NotifyHandler { + return &NotifyHandler{store: s, cache: c, matcher: m, renderer: r, router: rt} +} + +func (h *NotifyHandler) Handle(c *gin.Context) { + src := c.MustGet("source").(*model.Source) + + // 1. Read raw body + body, err := io.ReadAll(c.Request.Body) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": "read body failed"}) + return + } + + // 2. Parse message + p, err := parser.NewParser(src.ParseMode, src.ParsePattern) + if err != nil { + c.JSON(http.StatusUnprocessableEntity, gin.H{"error": "parser setup: " + err.Error()}) + return + } + msg, err := p.Parse(body) + if err != nil { + c.JSON(http.StatusUnprocessableEntity, gin.H{"error": "parse failed: " + err.Error()}) + return + } + + // 3. Match rule + rule, err := h.matcher.Match(c.Request.Context(), src.ID, msg.Event) + if err != nil { + // No matching rule → 200 with matched: false + c.JSON(http.StatusOK, gin.H{"matched": false}) + return + } + + // 4. Evaluate conditions + if rule.Conditions != nil { + var conds []model.Condition + if err := json.Unmarshal(*rule.Conditions, &conds); err == nil { + if !condition.Evaluate(conds, msg.Data) { + c.JSON(http.StatusOK, gin.H{ + "matched": true, + "filtered": true, + "reason": "condition not met", + }) + return + } + } + } + + // 5. Get template content + tmpl, err := h.store.GetTemplate(c.Request.Context(), rule.TemplateID) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": "template not found"}) + return + } + + // 6. Render template + content, err := h.renderer.Render(tmpl.Content, msg.Data) + if err != nil { + c.JSON(http.StatusUnprocessableEntity, gin.H{"error": "template render failed: " + err.Error()}) + return + } + + // 7. Route to channels + title := src.Name + ": " + msg.Event + channels := h.router.Route(c.Request.Context(), rule, title, content) + + // 8. Log message (best effort) + go func() { + payloadJSON, _ := json.Marshal(msg.Data) + for _, chName := range channels { + ml := &model.MessageLog{ + RuleID: rule.ID, + Source: src.Name, + Event: msg.Event, + Payload: payloadJSON, + Content: content, + Status: "pending", + } + _ = h.store.CreateMessageLog(c.Request.Context(), ml) + } + }() + + slog.Info("notification accepted", + "source", src.Name, + "event", msg.Event, + "channels", channels, + ) + + c.JSON(http.StatusOK, gin.H{ + "matched": true, + "channels": channels, + "accepted": true, + }) +} +``` + +Note: need imports for `cache` and `engine` — add: +```go +import ( + "aiaa-notification-service/internal/cache" +) +``` + +- [ ] **Step 2: Verify compilation** + +```bash +go build ./internal/handler/... +``` +Expected: no errors. + +- [ ] **Step 3: Commit** + +```bash +git add -A && git commit -m "feat: notify handler — parse, match, evaluate, render, route" +``` + +--- + +### Task 14: Management Handlers — Source, Template, Channel CRUD + +**Files:** +- Create: `internal/handler/source.go` +- Create: `internal/handler/template.go` +- Create: `internal/handler/channel.go` + +**Interfaces:** +- Consumes: `*store.Store`, `*cache.Cache` +- Produces: Source/Template/Channel CRUD endpoints + +- [ ] **Step 1: Write Source CRUD handler** + +File: `internal/handler/source.go` +```go +package handler + +import ( + "net/http" + "strconv" + + "aiaa-notification-service/internal/cache" + "aiaa-notification-service/internal/model" + "aiaa-notification-service/internal/store" + + "github.com/gin-gonic/gin" +) + +type SourceHandler struct { + store *store.Store + cache *cache.Cache +} + +func NewSourceHandler(s *store.Store, c *cache.Cache) *SourceHandler { + return &SourceHandler{store: s, cache: c} +} + +type createSourceReq struct { + Name string `json:"name" binding:"required"` + ParseMode string `json:"parse_mode"` + ParsePattern string `json:"parse_pattern"` + Status int `json:"status"` +} + +func (h *SourceHandler) Create(c *gin.Context) { + var req createSourceReq + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + if req.ParseMode == "" { + req.ParseMode = "json" + } + if req.Status == 0 { + req.Status = 1 + } + + src := &model.Source{ + Name: req.Name, + ParseMode: req.ParseMode, + ParsePattern: req.ParsePattern, + Status: req.Status, + } + if err := h.store.CreateSource(c.Request.Context(), src); err != nil { + c.JSON(http.StatusConflict, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusCreated, src) +} + +func (h *SourceHandler) List(c *gin.Context) { + sources, err := h.store.ListSources(c.Request.Context()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, sources) +} + +func (h *SourceHandler) Get(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + src, err := h.store.GetSource(c.Request.Context(), id) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "source not found"}) + return + } + c.JSON(http.StatusOK, src) +} + +func (h *SourceHandler) Update(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + var req createSourceReq + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + + src := &model.Source{ + Name: req.Name, + ParseMode: req.ParseMode, + ParsePattern: req.ParsePattern, + Status: req.Status, + } + if err := h.store.UpdateSource(c.Request.Context(), id, src); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + + // Invalidate cache + if h.cache != nil { + _ = h.cache.InvalidateBySource(c.Request.Context(), id) + } + + c.JSON(http.StatusOK, gin.H{"ok": true}) +} + +func (h *SourceHandler) Delete(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + if err := h.store.DeleteSource(c.Request.Context(), id); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) +} +``` + +- [ ] **Step 2: Write Template CRUD handler** + +File: `internal/handler/template.go` +```go +package handler + +import ( + "net/http" + "strconv" + + "aiaa-notification-service/internal/cache" + "aiaa-notification-service/internal/model" + "aiaa-notification-service/internal/store" + + "github.com/gin-gonic/gin" +) + +type TemplateHandler struct { + store *store.Store + cache *cache.Cache +} + +func NewTemplateHandler(s *store.Store, c *cache.Cache) *TemplateHandler { + return &TemplateHandler{store: s, cache: c} +} + +type createTemplateReq struct { + Name string `json:"name" binding:"required"` + Content string `json:"content" binding:"required"` +} + +func (h *TemplateHandler) Create(c *gin.Context) { + var req createTemplateReq + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + tmpl := &model.Template{Name: req.Name, Content: req.Content} + if err := h.store.CreateTemplate(c.Request.Context(), tmpl); err != nil { + c.JSON(http.StatusConflict, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusCreated, tmpl) +} + +func (h *TemplateHandler) List(c *gin.Context) { + templates, err := h.store.ListTemplates(c.Request.Context()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, templates) +} + +func (h *TemplateHandler) Get(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + tmpl, err := h.store.GetTemplate(c.Request.Context(), id) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "template not found"}) + return + } + c.JSON(http.StatusOK, tmpl) +} + +func (h *TemplateHandler) Update(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + var req createTemplateReq + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + tmpl := &model.Template{Name: req.Name, Content: req.Content} + if err := h.store.UpdateTemplate(c.Request.Context(), id, tmpl); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + + if h.cache != nil { + _ = h.cache.InvalidateByTemplate(c.Request.Context(), id) + } + + c.JSON(http.StatusOK, gin.H{"ok": true}) +} + +func (h *TemplateHandler) Delete(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + if err := h.store.DeleteTemplate(c.Request.Context(), id); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) +} +``` + +- [ ] **Step 3: Write Channel CRUD handler** + +File: `internal/handler/channel.go` +```go +package handler + +import ( + "encoding/json" + "net/http" + "strconv" + + "aiaa-notification-service/internal/cache" + "aiaa-notification-service/internal/model" + "aiaa-notification-service/internal/store" + + "github.com/gin-gonic/gin" +) + +type ChannelHandler struct { + store *store.Store + cache *cache.Cache +} + +func NewChannelHandler(s *store.Store, c *cache.Cache) *ChannelHandler { + return &ChannelHandler{store: s, cache: c} +} + +type createChannelReq struct { + Name string `json:"name" binding:"required"` + Type string `json:"type" binding:"required"` + Config json.RawMessage `json:"config" binding:"required"` + Status int `json:"status"` +} + +func (h *ChannelHandler) Create(c *gin.Context) { + var req createChannelReq + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + if req.Status == 0 { + req.Status = 1 + } + raw := json.RawMessage(req.Config) + ch := &model.Channel{Name: req.Name, Type: req.Type, Config: &raw, Status: req.Status} + if err := h.store.CreateChannel(c.Request.Context(), ch); err != nil { + c.JSON(http.StatusConflict, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusCreated, ch) +} + +func (h *ChannelHandler) List(c *gin.Context) { + channels, err := h.store.ListChannels(c.Request.Context()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, channels) +} + +func (h *ChannelHandler) Get(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + ch, err := h.store.GetChannel(c.Request.Context(), id) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "channel not found"}) + return + } + c.JSON(http.StatusOK, ch) +} + +func (h *ChannelHandler) Update(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + var req createChannelReq + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + raw := json.RawMessage(req.Config) + ch := &model.Channel{Name: req.Name, Type: req.Type, Config: &raw, Status: req.Status} + if err := h.store.UpdateChannel(c.Request.Context(), id, ch); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) +} + +func (h *ChannelHandler) Delete(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + if err := h.store.DeleteChannel(c.Request.Context(), id); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) +} +``` + +- [ ] **Step 4: Verify compilation** + +```bash +go build ./internal/handler/... +``` +Expected: no errors. + +- [ ] **Step 5: Commit** + +```bash +git add -A && git commit -m "feat: management handlers — source, template, channel CRUD" +``` + +--- + +### Task 15: Management Handlers — Rule CRUD + MessageLog + +**Files:** +- Create: `internal/handler/rule.go` +- Create: `internal/handler/message_log.go` + +**Interfaces:** +- Produces: Rule CRUD + enable/disable/channel-toggle endpoints +- Produces: GET message-logs with pagination + +- [ ] **Step 1: Write Rule CRUD handler** + +File: `internal/handler/rule.go` +```go +package handler + +import ( + "encoding/json" + "net/http" + "strconv" + + "aiaa-notification-service/internal/cache" + "aiaa-notification-service/internal/model" + "aiaa-notification-service/internal/store" + + "github.com/gin-gonic/gin" +) + +type RuleHandler struct { + store *store.Store + cache *cache.Cache +} + +func NewRuleHandler(s *store.Store, c *cache.Cache) *RuleHandler { + return &RuleHandler{store: s, cache: c} +} + +type createRuleReq struct { + SourceName string `json:"source_name" binding:"required"` + Event string `json:"event" binding:"required"` + TemplateName string `json:"template_name" binding:"required"` + Channels []string `json:"channels"` + Conditions []model.Condition `json:"conditions,omitempty"` + Enabled int `json:"enabled"` +} + +func (h *RuleHandler) Create(c *gin.Context) { + var req createRuleReq + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + if req.Enabled == 0 { + req.Enabled = 1 + } + + // Resolve names → IDs + src, err := h.store.GetSourceByName(c.Request.Context(), req.SourceName) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": "source not found: " + req.SourceName}) + return + } + tmpl, err := h.store.GetTemplateByName(c.Request.Context(), req.TemplateName) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": "template not found: " + req.TemplateName}) + return + } + channelIDs, err := resolveChannelNames(h.store, c, req.Channels) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + + var condsJSON *json.RawMessage + if len(req.Conditions) > 0 { + data, _ := json.Marshal(req.Conditions) + raw := json.RawMessage(data) + condsJSON = &raw + } + + rule := &model.Rule{ + SourceID: src.ID, + Event: req.Event, + TemplateID: tmpl.ID, + Conditions: condsJSON, + Enabled: req.Enabled, + } + + if err := h.store.CreateRule(c.Request.Context(), rule, channelIDs); err != nil { + c.JSON(http.StatusConflict, gin.H{"error": err.Error()}) + return + } + + c.JSON(http.StatusCreated, rule) +} + +func (h *RuleHandler) List(c *gin.Context) { + rules, err := h.store.ListRules(c.Request.Context()) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, rules) +} + +func (h *RuleHandler) Get(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + rule, err := h.store.GetRule(c.Request.Context(), id) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "rule not found"}) + return + } + c.JSON(http.StatusOK, rule) +} + +func (h *RuleHandler) Update(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + var req createRuleReq + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + + src, err := h.store.GetSourceByName(c.Request.Context(), req.SourceName) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": "source not found"}) + return + } + tmpl, err := h.store.GetTemplateByName(c.Request.Context(), req.TemplateName) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": "template not found"}) + return + } + channelIDs, err := resolveChannelNames(h.store, c, req.Channels) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + + var condsJSON *json.RawMessage + if len(req.Conditions) > 0 { + data, _ := json.Marshal(req.Conditions) + raw := json.RawMessage(data) + condsJSON = &raw + } + + rule := &model.Rule{ + SourceID: src.ID, + Event: req.Event, + TemplateID: tmpl.ID, + Conditions: condsJSON, + Enabled: req.Enabled, + } + + if err := h.store.UpdateRule(c.Request.Context(), id, rule, channelIDs); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + + // Invalidate cache + if h.cache != nil { + _ = h.cache.InvalidateRule(c.Request.Context(), src.ID, req.Event) + _ = h.cache.InvalidateChannels(c.Request.Context(), id) + } + + c.JSON(http.StatusOK, gin.H{"ok": true}) +} + +func (h *RuleHandler) Delete(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + if err := h.store.DeleteRule(c.Request.Context(), id); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) +} + +func (h *RuleHandler) Enable(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + if err := h.store.SetRuleEnabled(c.Request.Context(), id, true); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) +} + +func (h *RuleHandler) Disable(c *gin.Context) { + id, _ := strconv.Atoi(c.Param("id")) + if err := h.store.SetRuleEnabled(c.Request.Context(), id, false); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) +} + +func (h *RuleHandler) EnableChannel(c *gin.Context) { + ruleID, _ := strconv.Atoi(c.Param("id")) + channelID, _ := strconv.Atoi(c.Param("channel_id")) + if err := h.store.SetRuleChannelEnabled(c.Request.Context(), ruleID, channelID, true); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) +} + +func (h *RuleHandler) DisableChannel(c *gin.Context) { + ruleID, _ := strconv.Atoi(c.Param("id")) + channelID, _ := strconv.Atoi(c.Param("channel_id")) + if err := h.store.SetRuleChannelEnabled(c.Request.Context(), ruleID, channelID, false); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true}) +} + +func resolveChannelNames(s *store.Store, c *gin.Context, names []string) ([]int, error) { + var ids []int + for _, name := range names { + ch, err := s.GetChannelByName(c.Request.Context(), name) + if err != nil { + return nil, err + } + ids = append(ids, ch.ID) + } + return ids, nil +} +``` + +- [ ] **Step 2: Write MessageLog handler** + +File: `internal/handler/message_log.go` +```go +package handler + +import ( + "net/http" + + "aiaa-notification-service/internal/store" + + "github.com/gin-gonic/gin" +) + +type MessageLogHandler struct { + store *store.Store +} + +func NewMessageLogHandler(s *store.Store) *MessageLogHandler { + return &MessageLogHandler{store: s} +} + +func (h *MessageLogHandler) List(c *gin.Context) { + var filter store.MessageLogFilter + if err := c.ShouldBindQuery(&filter); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + + logs, total, err := h.store.ListMessageLogs(c.Request.Context(), filter) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + + c.JSON(http.StatusOK, gin.H{ + "data": logs, + "total": total, + "page": filter.Page, + }) +} +``` + +- [ ] **Step 3: Verify compilation** + +```bash +go build ./internal/handler/... +``` +Expected: no errors. + +- [ ] **Step 4: Commit** + +```bash +git add -A && git commit -m "feat: rule CRUD with enable/disable and message log query" +``` + +--- + +### Task 16: Retry Mechanism + +**Files:** +- Create: `internal/retry/retry.go` + +**Interfaces:** +- Produces: `retry.NewRetrier(maxRetries int, backoff []time.Duration) *Retrier`, `Retrier.Do(ctx context.Context, fn func() error) error` + +- [ ] **Step 1: Write retry module** + +File: `internal/retry/retry.go` +```go +package retry + +import ( + "context" + "fmt" + "log/slog" + "time" +) + +type Retrier struct { + maxRetries int + backoff []time.Duration +} + +func NewRetrier(maxRetries int, backoff []time.Duration) *Retrier { + return &Retrier{maxRetries: maxRetries, backoff: backoff} +} + +// DefaultRetrier returns a retrier with 3 attempts, exponential backoff: 1s, 5s, 30s. +func DefaultRetrier() *Retrier { + return NewRetrier(3, []time.Duration{1 * time.Second, 5 * time.Second, 30 * time.Second}) +} + +func (r *Retrier) Do(ctx context.Context, fn func() error) error { + var lastErr error + for attempt := 0; attempt <= r.maxRetries; attempt++ { + if attempt > 0 { + delay := r.backoff[attempt-1] + slog.Info("retrying", "attempt", attempt, "delay", delay) + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(delay): + } + } + + err := fn() + if err == nil { + return nil + } + lastErr = err + slog.Warn("attempt failed", "attempt", attempt, "error", err) + } + return fmt.Errorf("all %d attempts failed, last error: %w", r.maxRetries+1, lastErr) +} +``` + +- [ ] **Step 2: Verify compilation** + +```bash +go build ./internal/retry/... +``` +Expected: no errors. + +- [ ] **Step 3: Commit** + +```bash +git add -A && git commit -m "feat: retry mechanism with exponential backoff" +``` + +--- + +### Task 17: Main Assembly — Wire Everything Together + +**Files:** +- Modify: `cmd/server/main.go` +- Modify: `internal/handler/notify.go` (fix cache import) + +**Interfaces:** +- Produces: fully functional HTTP server + +- [ ] **Step 1: Write complete main.go** + +File: `cmd/server/main.go` +```go +package main + +import ( + "context" + "fmt" + "log/slog" + "net/http" + "os" + "os/signal" + "syscall" + "time" + + "aiaa-notification-service/internal/cache" + "aiaa-notification-service/internal/config" + "aiaa-notification-service/internal/engine" + "aiaa-notification-service/internal/handler" + "aiaa-notification-service/internal/store" + + "github.com/gin-gonic/gin" +) + +func main() { + logger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})) + slog.SetDefault(logger) + + // Load config + cfg, err := config.Load("config/config.yaml") + if err != nil { + slog.Error("failed to load config", "error", err) + os.Exit(1) + } + + // Connect MySQL + store, err := store.NewStore(cfg.Database) + if err != nil { + slog.Error("failed to connect to MySQL", "error", err) + os.Exit(1) + } + defer store.Close() + + // Connect Redis + var redisCache *cache.Cache + if cfg.Redis.Host != "" { + redisCache, err = cache.NewCache(cfg.Redis) + if err != nil { + slog.Warn("redis connection failed, running without cache", "error", err) + redisCache = nil + } + } else { + slog.Warn("redis not configured, running without cache") + } + if redisCache != nil { + defer redisCache.Close() + } + + // Build engine + matcher := engine.NewMatcher(store, redisCache) + renderer := engine.NewRenderer() + + // Build sender factory + senderFactory := func(channelType string) (adapter.ChannelSender, error) { + return adapter.NewSender(channelType, &cfg.SMTP) + } + router := engine.NewRouter(store, redisCache, senderFactory) + + // Build handlers + notifyH := handler.NewNotifyHandler(store, redisCache, matcher, renderer, router) + sourceH := handler.NewSourceHandler(store, redisCache) + templateH := handler.NewTemplateHandler(store, redisCache) + channelH := handler.NewChannelHandler(store, redisCache) + ruleH := handler.NewRuleHandler(store, redisCache) + msgLogH := handler.NewMessageLogHandler(store) + + // Setup Gin + gin.SetMode(gin.ReleaseMode) + r := gin.New() + r.Use(gin.Recovery()) + + // Health check + r.GET("/health", func(c *gin.Context) { + c.JSON(http.StatusOK, gin.H{"status": "ok"}) + }) + + api := r.Group("/api/v1") + + // Notify endpoint (source auth + rate limit) + notifyGroup := api.Group("/notify") + notifyGroup.Use(handler.SourceAuth(store)) + if redisCache != nil { + notifyGroup.Use(handler.RateLimit(redisCache, cfg.RateLimit.Default)) + } + notifyGroup.POST("", notifyH.Handle) + + // Management endpoints (admin auth) + admin := api.Group("") + admin.Use(handler.AdminAuth(cfg.Server.AdminKey)) + + // Sources + admin.POST("/sources", sourceH.Create) + admin.GET("/sources", sourceH.List) + admin.GET("/sources/:id", sourceH.Get) + admin.PUT("/sources/:id", sourceH.Update) + admin.DELETE("/sources/:id", sourceH.Delete) + + // Templates + admin.POST("/templates", templateH.Create) + admin.GET("/templates", templateH.List) + admin.GET("/templates/:id", templateH.Get) + admin.PUT("/templates/:id", templateH.Update) + admin.DELETE("/templates/:id", templateH.Delete) + + // Channels + admin.POST("/channels", channelH.Create) + admin.GET("/channels", channelH.List) + admin.GET("/channels/:id", channelH.Get) + admin.PUT("/channels/:id", channelH.Update) + admin.DELETE("/channels/:id", channelH.Delete) + + // Rules + admin.POST("/rules", ruleH.Create) + admin.GET("/rules", ruleH.List) + admin.GET("/rules/:id", ruleH.Get) + admin.PUT("/rules/:id", ruleH.Update) + admin.DELETE("/rules/:id", ruleH.Delete) + admin.PATCH("/rules/:id/enable", ruleH.Enable) + admin.PATCH("/rules/:id/disable", ruleH.Disable) + admin.PATCH("/rules/:id/channels/:channel_id/enable", ruleH.EnableChannel) + admin.PATCH("/rules/:id/channels/:channel_id/disable", ruleH.DisableChannel) + + // Message Logs + admin.GET("/message-logs", msgLogH.List) + + // Start server + srv := &http.Server{ + Addr: fmt.Sprintf(":%d", cfg.Server.Port), + Handler: r, + } + + go func() { + slog.Info("server starting", "port", cfg.Server.Port) + if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed { + slog.Error("server error", "error", err) + os.Exit(1) + } + }() + + // Graceful shutdown + quit := make(chan os.Signal, 1) + signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) + <-quit + slog.Info("shutting down...") + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if err := srv.Shutdown(ctx); err != nil { + slog.Error("forced shutdown", "error", err) + } + slog.Info("server stopped") +} +``` + +- [ ] **Step 2: Fix imports — add adapter import to main.go** + +Add to imports: +```go +import ( + "aiaa-notification-service/internal/adapter" +) +``` + +- [ ] **Step 3: Verify full build** + +```bash +go mod tidy +go build ./cmd/server/... +``` +Expected: builds successfully. Fix any import or type issues. + +- [ ] **Step 4: Commit** + +```bash +git add -A && git commit -m "feat: main assembly — wire all components, gin router, graceful shutdown" +``` + +--- + +### Task 18: Docker & Docker Compose + +**Files:** +- Create: `Dockerfile` +- Create: `docker-compose.yml` + +- [ ] **Step 1: Write Dockerfile** + +File: `Dockerfile` +```dockerfile +FROM golang:1.22-alpine AS builder + +WORKDIR /app +COPY go.mod go.sum ./ +RUN go mod download +COPY . . +RUN CGO_ENABLED=0 GOOS=linux go build -o /server ./cmd/server + +FROM alpine:3.20 + +RUN apk add --no-cache ca-certificates tzdata +COPY --from=builder /server /usr/local/bin/server +COPY config/config.yaml /etc/notification/config.yaml + +EXPOSE 8080 +ENTRYPOINT ["server"] +``` + +- [ ] **Step 2: Write docker-compose.yml** + +File: `docker-compose.yml` +```yaml +version: '3.8' + +services: + mysql: + image: mysql:8.0 + environment: + MYSQL_ROOT_PASSWORD: rootpass + MYSQL_DATABASE: notification + MYSQL_USER: notify + MYSQL_PASSWORD: notify + ports: + - "3306:3306" + volumes: + - mysql_data:/var/lib/mysql + healthcheck: + test: ["CMD", "mysqladmin", "ping", "-h", "localhost"] + interval: 5s + timeout: 3s + retries: 10 + + redis: + image: redis:7-alpine + ports: + - "6379:6379" + healthcheck: + test: ["CMD", "redis-cli", "ping"] + interval: 5s + timeout: 3s + retries: 5 + + api: + build: . + ports: + - "8080:8080" + environment: + - DB_PASSWORD=notify + - SMTP_PASSWORD= + depends_on: + mysql: + condition: service_healthy + redis: + condition: service_healthy + restart: unless-stopped + +volumes: + mysql_data: +``` + +- [ ] **Step 3: Update config for Docker** + +The `config/config.yaml` should use Docker service names for hosts: +```yaml +database: + host: "mysql" + port: 3306 + # ... + +redis: + host: "redis" + port: 6379 + # ... +``` + +Leave the existing config.yaml as-is (for local dev) and the docker-compose env vars will override as needed. + +- [ ] **Step 4: Commit** + +```bash +git add -A && git commit -m "feat: docker multi-stage build and docker-compose" +``` + +--- + +## Final Verification + +After all tasks complete, run the full test suite and build: + +```bash +# Tests +go test ./internal/... -v + +# Build +go build -o bin/server ./cmd/server + +# Start with docker-compose +docker-compose up -d + +# Run migrations (inside container or locally) +# Then test with curl: + +# Create a source +curl -X POST http://localhost:8080/api/v1/sources \ + -H "Authorization: Bearer admin-sk-change-me" \ + -H "Content-Type: application/json" \ + -d '{"name":"test-system","parse_mode":"json"}' + +# Create a template +curl -X POST http://localhost:8080/api/v1/templates \ + -H "Authorization: Bearer admin-sk-change-me" \ + -H "Content-Type: application/json" \ + -d '{"name":"test_tmpl","content":"🚀 {{.symbol}} {{.price}}"}' + +# Create a channel +curl -X POST http://localhost:8080/api/v1/channels \ + -H "Authorization: Bearer admin-sk-change-me" \ + -H "Content-Type: application/json" \ + -d '{"name":"test-bark","type":"bark","config":{"url":"https://api.day.app/test"}}' + +# Create a rule +curl -X POST http://localhost:8080/api/v1/rules \ + -H "Authorization: Bearer admin-sk-change-me" \ + -H "Content-Type: application/json" \ + -d '{"source_name":"test-system","event":"trade.open","template_name":"test_tmpl","channels":["test-bark"]}' + +# Send a notification (use the api_key from source creation response) +curl -X POST http://localhost:8080/api/v1/notify \ + -H "Authorization: Bearer " \ + -H "Content-Type: application/json" \ + -d '{"event":"trade.open","data":{"symbol":"BTC","price":65000}}' +``` + +--- diff --git a/docs/superpowers/specs/2025-06-27-notification-service-design.md b/docs/superpowers/specs/2025-06-27-notification-service-design.md new file mode 100644 index 0000000..5dfb38f --- /dev/null +++ b/docs/superpowers/specs/2025-06-27-notification-service-design.md @@ -0,0 +1,452 @@ +# Notification Service — 设计规格 + +## 概述 + +纯 Go 实现的通知服务。上游业务系统通过 webhook 推送消息,服务匹配规则、渲染模板、路由到多渠道发送。 + +**核心三层关系:** + +``` +source 1──N rule +template 1──N rule +rule N──M channel (via rule_channel) +``` + +## 技术栈 + +| 层 | 选型 | +|----|------| +| 语言 | Go 1.22+ | +| HTTP 框架 | gin-gonic/gin | +| 数据库驱动 | go-sql-driver/mysql | +| SQL 工具 | jmoiron/sqlx | +| 模板引擎 | text/template (标准库) | +| 缓存 | go-redis/redis | +| 配置 | viper | +| 日志 | slog (标准库) | +| 迁移 | golang-migrate | + +**外部依赖:** MySQL 8.0、Redis 7 + +**部署形态:** Docker 容器化 (docker-compose),最终镜像 ~12MB + +--- + +## 架构 + +``` +notification-service (Go) +│ +├── cmd/server/main.go ← 入口 +├── internal/ +│ ├── handler/ ← HTTP Handler (gin) +│ │ ├── notify.go ← POST /api/v1/notify (上游调用) +│ │ ├── source.go ← CRUD /api/v1/sources +│ │ ├── template.go ← CRUD /api/v1/templates +│ │ ├── channel.go ← CRUD /api/v1/channels +│ │ └── rule.go ← CRUD /api/v1/rules + 开关控制 +│ │ +│ ├── engine/ ← 核心引擎 +│ │ ├── matcher.go ← source + event → rule 匹配 +│ │ ├── renderer.go ← Go template 渲染 +│ │ └── router.go ← rule → channels 路由 +│ │ +│ ├── adapter/ ← 渠道适配器 +│ │ ├── adapter.go ← ChannelSender 接口 +│ │ ├── dingtalk.go ← 钉钉 webhook +│ │ ├── wecom.go ← 企业微信机器人 +│ │ ├── email.go ← SMTP 邮件 +│ │ └── bark.go ← Bark HTTP +│ │ +│ ├── model/ ← 数据模型 +│ ├── store/ ← MySQL 访问层 +│ └── cache/ ← Redis 缓存层 +│ +├── config/config.yaml +├── migrations/ +├── Dockerfile +├── docker-compose.yml +└── go.mod +``` + +--- + +## 数据模型 + +### source(来源系统) + +| 列 | 类型 | 说明 | +|----|------|------| +| id | INT PK | | +| name | VARCHAR(64) UNIQUE | trading-system | +| api_key | VARCHAR(128) UNIQUE | 通知 API 鉴权 token,创建 source 时自动生成 | +| parse_mode | VARCHAR(16) DEFAULT 'json' | json / regex / text(text=直接透传 body 到模板) | +| parse_pattern | VARCHAR(512) | regex 模式(含命名分组),parse_mode=regex 时必填 | +| status | TINYINT | 1=启用 0=禁用 | + +### template(模板) + +| 列 | 类型 | 说明 | +|----|------|------| +| id | INT PK | | +| name | VARCHAR(64) UNIQUE | trade_open | +| content | TEXT | Go template 语法,markdown/纯文本 | + +### channel(渠道配置) + +| 列 | 类型 | 说明 | +|----|------|------| +| id | INT PK | | +| name | VARCHAR(32) UNIQUE | dingtalk-prod | +| type | VARCHAR(32) | dingtalk / wecom / email / bark | +| config | JSON | 渠道特定配置 (webhook URL 等) | +| status | TINYINT | 1=启用 0=禁用 | + +### rule(规则 — 核心绑定) + +| 列 | 类型 | 说明 | +|----|------|------| +| id | INT PK | | +| source_id | INT FK→source | | +| event | VARCHAR(64) | trade.open / error / order.created | +| template_id | INT FK→template | | +| conditions | JSON | 过滤条件数组,null=不过滤。例: `[{"field":"symbol","op":"exists"}]` | +| enabled | TINYINT | 1=启用 0=禁用 | +| UNIQUE | (source_id, event) | 每个 source 下 event 唯一 | + +**条件操作符:** `eq` / `ne` / `gt` / `gte` / `lt` / `lte` / `exists` / `not_exists` / `contains` +多个条件 AND 关系,全部满足才转发。 + +### rule_channel(规则-渠道关联 — 独立开关) + +| 列 | 类型 | 说明 | +|----|------|------| +| id | INT PK | | +| rule_id | INT FK→rule | ON DELETE CASCADE | +| channel_id | INT FK→channel | | +| enabled | TINYINT | 每条规则下每个渠道可独立开关 | +| UNIQUE | (rule_id, channel_id) | | + +### message_log(消息记录 — 排查用) + +| 列 | 类型 | 说明 | +|----|------|------| +| id | BIGINT PK | | +| rule_id | INT | | +| channel_id | INT | | +| source | VARCHAR(64) | 冗余,方便查询 | +| event | VARCHAR(64) | 冗余,方便查询 | +| payload | JSON | 原始请求 data | +| content | TEXT | 渲染后内容 | +| status | ENUM | pending/success/failed/retrying | +| retry_count | INT | | +| response | TEXT | 渠道返回 | +| error_msg | TEXT | | +| created_at | DATETIME | | + +--- + +## API 设计 + +### 认证 + +- **通知接口**:`Authorization: Bearer ` — 通过 api_key 识别来源系统 +- 管理接口:`Authorization: Bearer ` — 全局管理 Key +- **Content-Type 不限** — source.parse_mode 决定如何解析 body: + - `json`: body 按 JSON 解析,`event` 取自 JSON 字段 + - `regex`: body 按正则解析,命名分组提取字段,`event` 需在 body 中显式提供 + - `text`: body 原样透传到模板(模板引用 `.Body`) + +### 通知接口 + +**JSON 模式(parse_mode=json):** +``` +POST /api/v1/notify +Authorization: Bearer +Content-Type: application/json + +{ + "event": "trade.open", + "data": { + "symbol": "BTC", + "price": 65000 + } +} +``` + +**Regex 模式(parse_mode=regex,parse_pattern 定义提取规则):** +``` +POST /api/v1/notify +Authorization: Bearer +Content-Type: text/plain + +BTC 开仓 价格:65000 方向:多 +``` + +**Text 模式(parse_mode=text,body 直接给模板):** +``` +POST /api/v1/notify +Authorization: Bearer +Content-Type: text/plain + +🚀 BTC 突破 68000,请注意风险 +``` + +**响应:** +``` +Response 200 (匹配且通过条件): +{ + "matched": true, + "channels": ["dingtalk", "email"], + "accepted": true +} + +Response 200 (匹配但条件不满足): +{ "matched": true, "filtered": true, "reason": "condition not met" } + +Response 200 (无匹配规则): +{ "matched": false } + +Response 401: +{ "error": "unauthorized" } + +Response 429: +{ "error": "rate_limit_exceeded", "message": "...", "retry_after": 1 } +``` + +### 管理 API + +``` +来源: +POST /api/v1/sources +GET /api/v1/sources +GET /api/v1/sources/:id +PUT /api/v1/sources/:id +DELETE /api/v1/sources/:id + +模板: +POST /api/v1/templates +GET /api/v1/templates +GET /api/v1/templates/:id +PUT /api/v1/templates/:id +DELETE /api/v1/templates/:id + +渠道: +POST /api/v1/channels +GET /api/v1/channels +GET /api/v1/channels/:id +PUT /api/v1/channels/:id +DELETE /api/v1/channels/:id + +规则: +POST /api/v1/rules +GET /api/v1/rules +GET /api/v1/rules/:id +PUT /api/v1/rules/:id +DELETE /api/v1/rules/:id +PATCH /api/v1/rules/:id/enable +PATCH /api/v1/rules/:id/disable +PATCH /api/v1/rules/:id/channels/:channel_id/enable +PATCH /api/v1/rules/:id/channels/:channel_id/disable + +消息记录: +GET /api/v1/message-logs?source=&event=&status=&page=1&page_size=20 +``` + +### 关键请求示例 + +**创建渠道:** +```json +POST /api/v1/channels +{ + "name": "dingtalk-prod", + "type": "dingtalk", + "config": { + "webhook_url": "https://oapi.dingtalk.com/robot/send?access_token=xxx", + "secret": "SEC..." + } +} +``` + +**创建规则(带条件过滤):** +```json +POST /api/v1/rules +{ + "source_name": "trading-system", + "event": "trade.open", + "template_name": "trade_open", + "channels": ["dingtalk", "email"], + "conditions": [ + {"field": "symbol", "op": "exists"}, + {"field": "price", "op": "gt", "value": "0"} + ] +} +``` + +**创建 source(Regex 解析模式):** +```json +POST /api/v1/sources +{ + "name": "legacy-monitor", + "parse_mode": "regex", + "parse_pattern": "(?P\\w+)\\s+(?P.+)" +} +``` + +--- + +## 核心执行流程 + +``` +POST /api/v1/notify + │ + ▼ +1. 鉴权 ── api_key 反查 source → source_id + parse_mode + parse_pattern + key 无效 → 401 + │ + ▼ +2. 消息解析 ── 根据 source.parse_mode: + json: body → JSON decode → event + data + regex: body → regex 匹配 → 命名分组 → data map + text: body → .Body 字段 + │ + ▼ +3. 规则匹配 ── 缓存: notify:rule:{source_id}:{event} + miss → SELECT * FROM rule WHERE source_id=? AND event=? AND enabled=1 + 无匹配 → 200 {matched: false} + │ + ▼ +4. 条件过滤 ── 检查 rule.conditions (如果非空) + 对每条 condition 求值 (eq/ne/gt/lt/exists/contains...) + 任一不满足 → 200 {matched: true, filtered: true} + │ + ▼ +5. 模板渲染 ── Go text/template 渲染 template.content + 解析后的 data + 失败 → 422 {error: "template render failed"} + │ + ▼ +6. 查找渠道 ── 缓存: notify:channels:{rule_id} + miss → SELECT * FROM rule_channel WHERE rule_id=? AND enabled=1 + 无可用渠道 → 200 {matched: true, sent: 0} + │ + ▼ +7. 并发发送 ── 每个渠道 goroutine 异步发送 + 记录 message_log + 失败 → 重试队列 (3次, 指数退避 1s/5s/30s) + │ + ▼ +8. 立即返回 200 {matched: true, channels: [...], accepted: true} +``` + +--- + +## 渠道适配器 + +```go +type ChannelSender interface { + Type() string + Send(title, content string, config json.RawMessage) error +} +``` + +| 渠道 | 协议 | 备注 | +|------|------|------| +| 钉钉 | POST webhook | 签名计算 + markdown 消息体 | +| 企业微信 | POST webhook | markdown/json 消息体 | +| 邮件 | SMTP | net/smtp 标准库,支持 TLS | +| Bark | POST HTTP | title + body POST 到设备 URL | + +新增渠道:实现 ChannelSender 接口 + 工厂注册一行。 + +--- + +## 缓存策略 + +| Key | 内容 | TTL | 失效 | +|-----|------|-----|------| +| `notify:rule:{source_id}:{event}` | rule_id + template_id + template_content | 5min | 规则/模板 CUD 时主动删除 | +| `notify:channels:{rule_id}` | [channel_id, type, config] | 5min | 规则-渠道变更时主动删除 | + +模式:Cache-Aside,管理 API 做 CUD 时主动失效,不依赖 TTL 被动过期。 + +--- + +## 限流 & 重试 + +**限流:** +- Per-source 限流,默认 100 req/s +- Redis 滑动窗口: `ratelimit:{source_id}:{window}` 1秒窗口 +- 超限返回 429 + Retry-After + +**重试:** +- 异步发送失败 → 内存重试队列 +- 最多 3 次,退避 1s → 5s → 30s +- 3 次全失败 → message_log.status = 'failed' +- 后续可扩展为 Redis 死信队列 + 手动重发 API + +--- + +## 配置 + +```yaml +server: + port: 8080 + admin_key: "admin-sk-xxx" + +database: + host: mysql + port: 3306 + user: notify + password: ${DB_PASSWORD} + database: notification + +redis: + host: redis + port: 6379 + +smtp: + host: smtp.example.com + port: 587 + user: notify@example.com + password: ${SMTP_PASSWORD} + +rate_limit: + default: 100 +``` + +--- + +## 部署 + +**docker-compose.yml:** +```yaml +services: + mysql: + image: mysql:8.0 + redis: + image: redis:7-alpine + api: + build: . + ports: ["8080:8080"] + depends_on: [mysql, redis] +``` + +**Dockerfile:** 多阶段构建 (golang:1.22-alpine → alpine:3.20),最终镜像 ~12MB。 + +--- + +## v1 范围 + +- [x] Source/Template/Channel/Rule CRUD 管理 API +- [x] POST /api/v1/notify 核心通知接口 +- [x] 消息解析:JSON / Regex / Text 三种模式 +- [x] 规则条件过滤(eq/ne/gt/lt/exists/not_exists/contains) +- [x] API Key 鉴权(source 用 source key,管理用 admin key) +- [x] 4 个渠道:钉钉、企业微信、邮件、Bark +- [x] Go template 渲染 +- [x] 规则级 + 规则-渠道级开关控制 +- [x] 异步发送 + 3 次指数退避重试 +- [x] Redis 缓存规则和渠道配置 +- [x] Per-source 限流 +- [x] message_log 消息记录 +- [x] Docker 部署 diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..4a71274 --- /dev/null +++ b/go.mod @@ -0,0 +1,20 @@ +module aiaa-notification-service + +go 1.26.2 + +require github.com/spf13/viper v1.21.0 + +require ( + github.com/fsnotify/fsnotify v1.9.0 // indirect + github.com/go-viper/mapstructure/v2 v2.4.0 // indirect + github.com/pelletier/go-toml/v2 v2.2.4 // indirect + github.com/sagikazarmark/locafero v0.11.0 // indirect + github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 // indirect + github.com/spf13/afero v1.15.0 // indirect + github.com/spf13/cast v1.10.0 // indirect + github.com/spf13/pflag v1.0.10 // indirect + github.com/subosito/gotenv v1.6.0 // indirect + go.yaml.in/yaml/v3 v3.0.4 // indirect + golang.org/x/sys v0.29.0 // indirect + golang.org/x/text v0.28.0 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..81f59c4 --- /dev/null +++ b/go.sum @@ -0,0 +1,47 @@ +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/frankban/quicktest v1.14.6 h1:7Xjx+VpznH+oBnejlPUj8oUpdxnVs4f8XU8WnHkI4W8= +github.com/frankban/quicktest v1.14.6/go.mod h1:4ptaffx2x8+WTWXmUCuVU6aPUX1/Mz7zb5vbUoiM6w0= +github.com/fsnotify/fsnotify v1.9.0 h1:2Ml+OJNzbYCTzsxtv8vKSFD9PbJjmhYF14k/jKC7S9k= +github.com/fsnotify/fsnotify v1.9.0/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0= +github.com/go-viper/mapstructure/v2 v2.4.0 h1:EBsztssimR/CONLSZZ04E8qAkxNYq4Qp9LvH92wZUgs= +github.com/go-viper/mapstructure/v2 v2.4.0/go.mod h1:oJDH3BJKyqBA2TXFhDsKDGDTlndYOZ6rGS0BRZIxGhM= +github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= +github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0t5Ec4= +github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/rogpeppe/go-internal v1.9.0 h1:73kH8U+JUqXU8lRuOHeVHaa/SZPifC7BkcraZVejAe8= +github.com/rogpeppe/go-internal v1.9.0/go.mod h1:WtVeX8xhTBvf0smdhujwtBcq4Qrzq/fJaraNFVN+nFs= +github.com/sagikazarmark/locafero v0.11.0 h1:1iurJgmM9G3PA/I+wWYIOw/5SyBtxapeHDcg+AAIFXc= +github.com/sagikazarmark/locafero v0.11.0/go.mod h1:nVIGvgyzw595SUSUE6tvCp3YYTeHs15MvlmU87WwIik= +github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 h1:+jumHNA0Wrelhe64i8F6HNlS8pkoyMv5sreGx2Ry5Rw= +github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8/go.mod h1:3n1Cwaq1E1/1lhQhtRK2ts/ZwZEhjcQeJQ1RuC6Q/8U= +github.com/spf13/afero v1.15.0 h1:b/YBCLWAJdFWJTN9cLhiXXcD7mzKn9Dm86dNnfyQw1I= +github.com/spf13/afero v1.15.0/go.mod h1:NC2ByUVxtQs4b3sIUphxK0NioZnmxgyCrfzeuq8lxMg= +github.com/spf13/cast v1.10.0 h1:h2x0u2shc1QuLHfxi+cTJvs30+ZAHOGRic8uyGTDWxY= +github.com/spf13/cast v1.10.0/go.mod h1:jNfB8QC9IA6ZuY2ZjDp0KtFO2LZZlg4S/7bzP6qqeHo= +github.com/spf13/pflag v1.0.10 h1:4EBh2KAYBwaONj6b2Ye1GiHfwjqyROoF4RwYO+vPwFk= +github.com/spf13/pflag v1.0.10/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg= +github.com/spf13/viper v1.21.0 h1:x5S+0EU27Lbphp4UKm1C+1oQO+rKx36vfCoaVebLFSU= +github.com/spf13/viper v1.21.0/go.mod h1:P0lhsswPGWD/1lZJ9ny3fYnVqxiegrlNrEmgLjbTCAY= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8= +github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU= +go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= +go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= +golang.org/x/sys v0.29.0 h1:TPYlXGxvx1MGTn2GiZDhnjPA9wZzZeGKHHmKhHYvgaU= +golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/text v0.28.0 h1:rhazDwis8INMIwQ4tpjLDzUhx6RlXqZNPEM0huQojng= +golang.org/x/text v0.28.0/go.mod h1:U8nCwOR8jO/marOQ0QbDiOngZVEBB7MAiitBuMjXiNU= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15 h1:YR8cESwS4TdDjEe65xsg0ogRM/Nc3DYOhEAlW+xobZo= +gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/internal/config/config.go b/internal/config/config.go new file mode 100644 index 0000000..2832b5c --- /dev/null +++ b/internal/config/config.go @@ -0,0 +1,100 @@ +package config + +import ( + "fmt" + "os" + "strings" + + "github.com/spf13/viper" +) + +type Config struct { + Server ServerConfig `mapstructure:"server"` + Database DatabaseConfig `mapstructure:"database"` + Redis RedisConfig `mapstructure:"redis"` + SMTP SMTPConfig `mapstructure:"smtp"` + RateLimit RateLimitConfig `mapstructure:"rate_limit"` +} + +type ServerConfig struct { + Port int `mapstructure:"port"` + AdminKey string `mapstructure:"admin_key"` +} + +type DatabaseConfig struct { + Host string `mapstructure:"host"` + Port int `mapstructure:"port"` + User string `mapstructure:"user"` + Password string `mapstructure:"password"` + Database string `mapstructure:"database"` +} + +func (d DatabaseConfig) DSN() string { + return fmt.Sprintf("%s:%s@tcp(%s:%d)/%s?charset=utf8mb4&parseTime=true&loc=Local", + d.User, d.Password, d.Host, d.Port, d.Database) +} + +type RedisConfig struct { + Host string `mapstructure:"host"` + Port int `mapstructure:"port"` + Password string `mapstructure:"password"` + DB int `mapstructure:"db"` +} + +func (r RedisConfig) Addr() string { + return fmt.Sprintf("%s:%d", r.Host, r.Port) +} + +type SMTPConfig struct { + Host string `mapstructure:"host"` + Port int `mapstructure:"port"` + User string `mapstructure:"user"` + Password string `mapstructure:"password"` + From string `mapstructure:"from"` +} + +type RateLimitConfig struct { + Default int `mapstructure:"default"` +} + +func Load(path string) (*Config, error) { + v := viper.New() + v.SetConfigFile(path) + v.SetEnvPrefix("NOTIFY") + v.SetEnvKeyReplacer(strings.NewReplacer(".", "_")) + v.AutomaticEnv() + + if err := v.ReadInConfig(); err != nil { + return nil, fmt.Errorf("read config: %w", err) + } + + // Expand ${ENV_VAR} placeholders in config values + for _, key := range v.AllKeys() { + val := v.GetString(key) + if strings.Contains(val, "${") { + expanded := expandEnv(val) + v.Set(key, expanded) + } + } + + var cfg Config + if err := v.Unmarshal(&cfg); err != nil { + return nil, fmt.Errorf("unmarshal config: %w", err) + } + return &cfg, nil +} + +func expandEnv(s string) string { + return os.Expand(s, func(key string) string { + // support ${VAR:-default} + if i := strings.Index(key, ":-"); i >= 0 { + name := key[:i] + def := key[i+2:] + if v, ok := os.LookupEnv(name); ok { + return v + } + return def + } + return os.Getenv(key) + }) +}