Skip to content

Commit 105c111

Browse files
committed
feat(log): 全栈结构化日志治理
横切层 - middleware/request_logger.go 新增 HTTP 访问日志,抽取/生成 X-Request-ID 注入响应头与 std ctx,按状态分级(健康检查降为 Debug) - middleware/recovery.go 新增 panic 兜底,输出 panic_recovered + 完整 stack,返 500 + request_id - middleware/auth.go 鉴权失败带 reason 的 *_validation_failed Warn, 成功后用 user_id/group_id/api_key_id 重派生 logger 注入下游 - server/router.go server.go 启动顺序:Recovery → RequestLogger → I18n 业务路径 - plugin/forwarder.go forward_request_start 降 Debug、单请求 1 行 INFO 策略;空 attrs / nil error 不输出;429 + Retry-After 取代静默 503 - plugin/request.go body 读失败 / 平台插件未加载 / 路由缺失补错误日志 - plugin/manager_runtime.go plugin_load/start_timeout 等 + 字段名常量化 - plugin/host_service.go host_forward_* 系列结构化日志 - plugin/outcome.go 新增 openAIRateLimitError helper 服务于 forward 限流 - auth/{apikey,crypto}.go 缓存命中/解密失败补点 - billing/recorder.go 事件名规范化 - scheduler/{state,selection}.go 状态翻转保留 Info、降级 Warn - routing/selector.go routing_match Debug / routing_no_match Warn 管理面 12 个 app 包审计日志 - {auth,user,apikey,dashboard,account,group,subscription,pluginadmin, settings,usage,openclaw,proxy}/service.go 关键操作(创建/删除/改密/ 改额度/状态翻转)补 Info 审计;失败/拒绝补 Warn/Error;密码/token/ 签名/key 一律不打印 基础设施 - infra/store/ent_logger.go 新增 ent 默认 logger 桥接到 slog.Debug - infra/mailer/mailer.go SMTP 三层日志,邮箱用 hash 脱敏 - setup/upgrade/bootstrap 一次性事件 + 失败 Error - cmd/server/main.go server_listening / config_loaded(DSN 脱敏) config.yaml.example 增加 log: 块说明 level/format 取值
1 parent ef8df3f commit 105c111

36 files changed

Lines changed: 1650 additions & 225 deletions

File tree

backend/cmd/server/main.go

Lines changed: 52 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import (
2626
"github.com/DouDOU-start/airgate-core/internal/bootstrap"
2727
"github.com/DouDOU-start/airgate-core/internal/config"
2828
"github.com/DouDOU-start/airgate-core/internal/i18n"
29+
"github.com/DouDOU-start/airgate-core/internal/infra/store"
2930
"github.com/DouDOU-start/airgate-core/internal/server"
3031
"github.com/DouDOU-start/airgate-core/internal/setup"
3132
"github.com/DouDOU-start/airgate-core/internal/version"
@@ -70,14 +71,16 @@ func main() {
7071
}
7172

7273
// 加载配置
73-
cfg, err := config.Load(config.ConfigPath())
74+
cfgPath := config.ConfigPath()
75+
cfg, err := config.Load(cfgPath)
7476
if err != nil {
75-
slog.Error("加载配置失败", "error", err)
77+
slog.Error("config_load_failed", "path", cfgPath, "error", err)
7678
os.Exit(1)
7779
}
7880

7981
// 用配置值重新初始化日志(应用配置文件中的 level/format)
8082
sdk.InitLogger("core", cfg.Log.Level, cfg.Log.Format)
83+
slog.Info("config_loaded", "path", cfgPath, "log_level", cfg.Log.Level, "log_format", cfg.Log.Format)
8184

8285
// 启动正常服务
8386
startMainServer(cfg)
@@ -135,12 +138,22 @@ func startSetupServer() {
135138

136139
// startMainServer 启动主服务器
137140
func startMainServer(cfg *config.Config) {
141+
bootStart := time.Now()
142+
138143
// 初始化数据库连接(Ent Client)
139-
drv, err := sql.Open(dialect.Postgres, cfg.Database.DSN())
144+
dsn := cfg.Database.DSN()
145+
drv, err := sql.Open(dialect.Postgres, dsn)
140146
if err != nil {
141-
slog.Error("打开数据库失败", "error", err)
147+
slog.Error("db_open_failed", "dsn", store.RedactDSN(dsn), sdk.LogFieldError, err)
142148
os.Exit(1)
143149
}
150+
slog.Info("db_connected",
151+
"driver", "postgres",
152+
"host", cfg.Database.Host,
153+
"port", cfg.Database.Port,
154+
"db", cfg.Database.DBName,
155+
"dsn", store.RedactDSN(dsn))
156+
144157
// 配置连接池:不限制时 Go 会无限开连接,高并发下 Postgres "too many clients already"
145158
maxOpen := cfg.Database.MaxOpenConns
146159
if maxOpen <= 0 {
@@ -157,19 +170,33 @@ func startMainServer(cfg *config.Config) {
157170
drv.DB().SetMaxOpenConns(maxOpen)
158171
drv.DB().SetMaxIdleConns(maxIdle)
159172
drv.DB().SetConnMaxLifetime(time.Duration(lifeMin) * time.Minute)
160-
slog.Info("数据库连接池已配置",
173+
slog.Info("db_pool_configured",
161174
"max_open", maxOpen, "max_idle", maxIdle, "lifetime_min", lifeMin)
162175

163-
db := ent.NewClient(ent.Driver(drv))
176+
pingCtx, pingCancel := context.WithTimeout(context.Background(), 5*time.Second)
177+
if err := drv.DB().PingContext(pingCtx); err != nil {
178+
pingCancel()
179+
slog.Error("db_ping_failed", sdk.LogFieldError, err)
180+
os.Exit(1)
181+
}
182+
pingCancel()
183+
184+
// 注入 slog 桥接,让 ent 内部 debug/error 日志走结构化通道
185+
db := ent.NewClient(ent.Driver(drv), store.EntSlogLogger())
186+
slog.Info("ent_client_initialized",
187+
"driver", "postgres",
188+
"max_open_conns", maxOpen,
189+
"max_idle_conns", maxIdle,
190+
"lifetime_min", lifeMin)
164191
defer func() {
165192
if err := db.Close(); err != nil {
166-
slog.Warn("关闭数据库连接失败", "error", err)
193+
slog.Warn("db_close_failed", sdk.LogFieldError, err)
167194
}
168195
}()
169196

170197
// 启动时执行非破坏性迁移,补齐缺失表和字段,避免升级后因 schema 落后导致接口报错。
171198
if err := db.Schema.Create(context.Background(), migrate.WithDropIndex(false), migrate.WithDropColumn(false)); err != nil {
172-
slog.Error("执行数据库迁移失败", "error", err)
199+
slog.Error("db_migration_failed", sdk.LogFieldError, err)
173200
os.Exit(1)
174201
}
175202

@@ -182,12 +209,28 @@ func startMainServer(cfg *config.Config) {
182209
Password: cfg.Redis.Password,
183210
DB: cfg.Redis.DB,
184211
})
212+
redisCtx, redisCancel := context.WithTimeout(context.Background(), 5*time.Second)
213+
if err := rdb.Ping(redisCtx).Err(); err != nil {
214+
redisCancel()
215+
slog.Error("redis_ping_failed",
216+
"host", cfg.Redis.Host,
217+
"port", cfg.Redis.Port,
218+
sdk.LogFieldError, err)
219+
os.Exit(1)
220+
}
221+
redisCancel()
222+
slog.Info("redis_connected",
223+
"host", cfg.Redis.Host,
224+
"port", cfg.Redis.Port,
225+
"db", cfg.Redis.DB)
185226
defer func() {
186227
if err := rdb.Close(); err != nil {
187-
slog.Warn("关闭 Redis 连接失败", "error", err)
228+
slog.Warn("redis_close_failed", sdk.LogFieldError, err)
188229
}
189230
}()
190231

232+
slog.Info("bootstrap_completed", "duration_ms", time.Since(bootStart).Milliseconds())
233+
191234
// 创建并启动 HTTP 服务器
192235
srv := server.NewServer(cfg, db, rdb)
193236

backend/config.yaml.example

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,12 @@ security:
2525
api_key_secret: "" # 可选,API Key 加密密钥(hex 编码,≥64 字符)
2626
# 未配置时使用内置默认值;或设置环境变量 API_KEY_SECRET
2727

28+
log:
29+
level: "info" # debug / info / warn / error;core 设置会自动透传给所有插件
30+
# 或设置环境变量 LOG_LEVEL
31+
format: "text" # text(人看)/ json(对接 Loki / ELK / Datadog)
32+
# 或设置环境变量 LOG_FORMAT
33+
2834
plugins:
2935
dir: data/plugins
3036
dev: []

backend/internal/app/account/service.go

Lines changed: 117 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -92,14 +92,38 @@ func (s *Service) List(ctx context.Context, filter ListFilter) (ListResult, erro
9292
}
9393

9494
ids := make([]int, 0, len(accounts))
95+
openaiIDs := make([]int, 0, len(accounts))
9596
for _, item := range accounts {
9697
ids = append(ids, item.ID)
98+
// 生图统计仅 OpenAI 平台账号需要:其它平台没有 image endpoint,跑 SQL 也是 0 行白浪费。
99+
if item.Platform == "openai" {
100+
openaiIDs = append(openaiIDs, item.ID)
101+
}
97102
}
98103
counts := s.concurrency.GetCurrentCounts(ctx, ids)
99104
for index := range accounts {
100105
accounts[index].CurrentConcurrency = counts[accounts[index].ID]
101106
}
102107

108+
// 生图请求计数:今日 + 累计。BatchImageStats 失败不阻断主响应(运维路径优先稳定)。
109+
if len(openaiIDs) > 0 {
110+
todayStart := timezone.StartOfDay(s.now().In(time.Local))
111+
if imageStats, err := s.repo.BatchImageStats(ctx, openaiIDs, todayStart); err == nil {
112+
for index := range accounts {
113+
if accounts[index].Platform != "openai" {
114+
continue
115+
}
116+
if entry, ok := imageStats[accounts[index].ID]; ok {
117+
stats := entry
118+
accounts[index].ImageStats = &stats
119+
} else {
120+
// 没记录:显式给个零值结构,让前端拿到 today=0/total=0 而不是 nil(区分"没数据"和"非 openai")
121+
accounts[index].ImageStats = &AccountImageStats{}
122+
}
123+
}
124+
}
125+
}
126+
103127
return ListResult{
104128
List: accounts,
105129
Total: total,
@@ -110,10 +134,22 @@ func (s *Service) List(ctx context.Context, filter ListFilter) (ListResult, erro
110134

111135
// Create 创建账号。
112136
func (s *Service) Create(ctx context.Context, input CreateInput) (Account, error) {
137+
logger := sdk.LoggerFromContext(ctx)
113138
account, err := s.repo.Create(ctx, input)
114-
if err == nil {
115-
s.InvalidateUsageCache("") // 新账号创建后清除用量缓存
116-
}
139+
if err != nil {
140+
logger.Error("account_credential_persist_failed",
141+
sdk.LogFieldPlatform, input.Platform,
142+
"type", input.Type,
143+
"name", input.Name,
144+
sdk.LogFieldError, err)
145+
return account, err
146+
}
147+
logger.Info("account_created",
148+
sdk.LogFieldAccountID, account.ID,
149+
sdk.LogFieldPlatform, account.Platform,
150+
"type", account.Type,
151+
"name", account.Name)
152+
s.InvalidateUsageCache("") // 新账号创建后清除用量缓存
117153
return account, err
118154
}
119155

@@ -144,15 +180,39 @@ func (s *Service) Import(ctx context.Context, items []CreateInput) ImportSummary
144180

145181
// Update 更新账号。
146182
func (s *Service) Update(ctx context.Context, id int, input UpdateInput) (Account, error) {
147-
return s.repo.Update(ctx, id, input)
183+
logger := sdk.LoggerFromContext(ctx)
184+
updated, err := s.repo.Update(ctx, id, input)
185+
if err != nil {
186+
logger.Error("account_credential_persist_failed",
187+
sdk.LogFieldAccountID, id,
188+
sdk.LogFieldError, err)
189+
return updated, err
190+
}
191+
switch {
192+
case input.State != nil:
193+
logger.Info("account_status_changed",
194+
sdk.LogFieldAccountID, id,
195+
"state", *input.State)
196+
case input.MaxConcurrency != nil || input.RateMultiplier != nil:
197+
logger.Info("account_quota_updated",
198+
sdk.LogFieldAccountID, id)
199+
}
200+
return updated, err
148201
}
149202

150203
// Delete 删除账号。
151204
func (s *Service) Delete(ctx context.Context, id int) error {
205+
logger := sdk.LoggerFromContext(ctx)
152206
err := s.repo.Delete(ctx, id)
153-
if err == nil {
154-
s.InvalidateUsageCache("")
207+
if err != nil {
208+
logger.Error("account_credential_persist_failed",
209+
sdk.LogFieldAccountID, id,
210+
"op", "delete",
211+
sdk.LogFieldError, err)
212+
return err
155213
}
214+
logger.Info("account_deleted", sdk.LogFieldAccountID, id)
215+
s.InvalidateUsageCache("")
156216
return err
157217
}
158218

@@ -212,8 +272,12 @@ func (r *BulkResult) appendFailure(id int, err error) {
212272
// ToggleScheduling 快速切换账号调度状态。active ↔ disabled。
213273
// 其它中间态(rate_limited / degraded)一律视为"非 disabled",切换后目标 = disabled。
214274
func (s *Service) ToggleScheduling(ctx context.Context, id int) (ToggleResult, error) {
275+
logger := sdk.LoggerFromContext(ctx)
215276
item, err := s.repo.FindByID(ctx, id, LoadOptions{})
216277
if err != nil {
278+
logger.Error("account_lookup_failed",
279+
sdk.LogFieldAccountID, id,
280+
sdk.LogFieldError, err)
217281
return ToggleResult{}, err
218282
}
219283

@@ -224,20 +288,35 @@ func (s *Service) ToggleScheduling(ctx context.Context, id int) (ToggleResult, e
224288

225289
updated, err := s.repo.Update(ctx, id, UpdateInput{State: &newState})
226290
if err != nil {
291+
logger.Error("account_credential_persist_failed",
292+
sdk.LogFieldAccountID, id,
293+
"op", "toggle_scheduling",
294+
sdk.LogFieldError, err)
227295
return ToggleResult{}, err
228296
}
297+
logger.Info("account_status_changed",
298+
sdk.LogFieldAccountID, id,
299+
"state", updated.State)
229300
return ToggleResult{ID: updated.ID, State: updated.State}, nil
230301
}
231302

232303
// PrepareConnectivityTest 准备账号连通性测试。
233304
func (s *Service) PrepareConnectivityTest(ctx context.Context, id int, modelID string) (*ConnectivityTest, error) {
305+
logger := sdk.LoggerFromContext(ctx)
234306
item, err := s.repo.FindByID(ctx, id, LoadOptions{WithProxy: true})
235307
if err != nil {
308+
logger.Error("account_lookup_failed",
309+
sdk.LogFieldAccountID, id,
310+
sdk.LogFieldError, err)
236311
return nil, err
237312
}
238313

239314
inst := s.plugins.GetPluginByPlatform(item.Platform)
240315
if inst == nil || inst.Gateway == nil {
316+
logger.Warn("account_credential_validation_failed",
317+
sdk.LogFieldAccountID, id,
318+
sdk.LogFieldPlatform, item.Platform,
319+
sdk.LogFieldReason, "plugin_not_found")
241320
return nil, ErrPluginNotFound
242321
}
243322

@@ -808,13 +887,21 @@ func (s *Service) GetCredentialsSchema(platform string) CredentialSchema {
808887

809888
// RefreshQuota 刷新账号额度。
810889
func (s *Service) RefreshQuota(ctx context.Context, id int) (QuotaRefreshResult, error) {
890+
logger := sdk.LoggerFromContext(ctx)
811891
item, err := s.repo.FindByID(ctx, id, LoadOptions{})
812892
if err != nil {
893+
logger.Error("account_lookup_failed",
894+
sdk.LogFieldAccountID, id,
895+
sdk.LogFieldError, err)
813896
return QuotaRefreshResult{}, err
814897
}
815898

816899
inst := s.plugins.GetPluginByPlatform(item.Platform)
817900
if inst == nil || inst.Gateway == nil {
901+
logger.Warn("account_credential_validation_failed",
902+
sdk.LogFieldAccountID, id,
903+
sdk.LogFieldPlatform, item.Platform,
904+
sdk.LogFieldReason, "quota_refresh_unsupported")
818905
return QuotaRefreshResult{}, ErrQuotaRefreshUnsupported
819906
}
820907

@@ -825,8 +912,16 @@ func (s *Service) RefreshQuota(ctx context.Context, id int) (QuotaRefreshResult,
825912
if err != nil {
826913
// 识别插件返回的 reauth_required 前缀(字符串识别,gRPC 不透传 sentinel error)。
827914
if strings.Contains(err.Error(), reauthRequiredPrefix) {
915+
logger.Warn("account_credential_validation_failed",
916+
sdk.LogFieldAccountID, id,
917+
sdk.LogFieldPlatform, item.Platform,
918+
sdk.LogFieldReason, "reauth_required")
828919
return QuotaRefreshResult{}, ErrReauthRequired
829920
}
921+
logger.Error("account_credential_validation_failed",
922+
sdk.LogFieldAccountID, id,
923+
sdk.LogFieldPlatform, item.Platform,
924+
sdk.LogFieldError, err)
830925
return QuotaRefreshResult{}, fmt.Errorf("刷新额度失败: %w", err)
831926
}
832927

@@ -853,6 +948,10 @@ func (s *Service) RefreshQuota(ctx context.Context, id int) (QuotaRefreshResult,
853948
}
854949
if updated {
855950
if err := s.repo.SaveCredentials(ctx, id, credentials); err != nil {
951+
logger.Error("account_credential_persist_failed",
952+
sdk.LogFieldAccountID, id,
953+
"op", "save_credentials",
954+
sdk.LogFieldError, err)
856955
return QuotaRefreshResult{}, err
857956
}
858957
}
@@ -885,17 +984,23 @@ func (s *Service) triggerUsageProbe(ctx context.Context, inst *plugin.PluginInst
885984
defer cancel()
886985
status, _, _, err := inst.Gateway.HandleHTTPRequest(probeCtx, "POST", "usage/probe", "", nil, reqBody)
887986
if err != nil || status != http.StatusOK {
888-
slog.Debug("usage/probe 探测失败(降级:等下一轮 5 分钟缓存过期重试)",
889-
"account_id", id, "status", status, "error", err)
987+
slog.Debug("account_usage_probe_failed",
988+
sdk.LogFieldAccountID, id,
989+
sdk.LogFieldStatus, status,
990+
sdk.LogFieldError, err)
890991
}
891992
// 清掉本进程的 usage 5 分钟缓存,让下一次 GetAccountUsage 重新从插件拉窗口数据。
892993
s.InvalidateUsageCache("")
893994
}
894995

895996
// GetStats 获取单个账号统计。
896997
func (s *Service) GetStats(ctx context.Context, id int, query StatsQuery) (StatsResult, error) {
998+
logger := sdk.LoggerFromContext(ctx)
897999
item, err := s.repo.FindByID(ctx, id, LoadOptions{})
8981000
if err != nil {
1001+
logger.Error("account_lookup_failed",
1002+
sdk.LogFieldAccountID, id,
1003+
sdk.LogFieldError, err)
8991004
return StatsResult{}, err
9001005
}
9011006

@@ -908,6 +1013,10 @@ func (s *Service) GetStats(ctx context.Context, id int, query StatsQuery) (Stats
9081013

9091014
logs, err := s.repo.FindUsageLogs(ctx, id, startDate, endDate)
9101015
if err != nil {
1016+
logger.Error("account_lookup_failed",
1017+
sdk.LogFieldAccountID, id,
1018+
"op", "find_usage_logs",
1019+
sdk.LogFieldError, err)
9111020
return StatsResult{}, err
9121021
}
9131022

0 commit comments

Comments
 (0)