32 changed files with 1631 additions and 39 deletions
@ -0,0 +1,80 @@ |
|||||
|
package comm |
||||
|
|
||||
|
import ( |
||||
|
"crypto/hmac" |
||||
|
"crypto/sha256" |
||||
|
"encoding/hex" |
||||
|
"errors" |
||||
|
"fmt" |
||||
|
"strconv" |
||||
|
"time" |
||||
|
) |
||||
|
|
||||
|
// console → 业务服务的「运维调用」鉴权。
|
||||
|
//
|
||||
|
// 背景:console 是单例服务(lego/base/single),不接 ETCD/rpcx 集群,无法直接 RPC 业务服务;
|
||||
|
// 它只能走业务网关的公网 HTTP 入口。而这类接口(如"重载业务配置")绝不能裸奔在公网上,
|
||||
|
// 故用 console 与各业务服务**本就必须一致**的 ${FIELD_ENCRYPT_KEY} 做 HMAC 共享密钥签名——
|
||||
|
// 不新增需要两边同步的配置项(多一把密钥就多一次"两边不一致"的事故)。
|
||||
|
//
|
||||
|
// 签名串 = HMAC-SHA256(key, "<ts>.<msgName>"),绑定了时间与目标接口;
|
||||
|
// ts 偏差超过 consoleSignSkew 即拒绝,限制重放窗口。
|
||||
|
|
||||
|
const ( |
||||
|
ConsoleTsHeader = "X-Console-Ts" // 秒级 unix 时间戳
|
||||
|
ConsoleSignHeader = "X-Console-Sign" // hex(HMAC-SHA256)
|
||||
|
ConsoleArgHeader = "X-Console-Arg" // 接口参数(如模块名)。已并入签名,不可篡改
|
||||
|
|
||||
|
// SessionMeta_ConsoleArg 网关把 ConsoleArgHeader 搬进 args.Meta 用的键。
|
||||
|
// 业务侧 handler 通过 session.GetMateToString 读取——这样参数不必新增 pb 字段。
|
||||
|
SessionMeta_ConsoleArg = "console_arg" |
||||
|
|
||||
|
// consoleSignSkew 允许的时间偏差。窗口内同一签名可被重放,但影响仅限于"多触发一次重载/重置",
|
||||
|
// 二者都幂等;再收紧会被两机时钟漂移误伤。
|
||||
|
consoleSignSkew = 5 * time.Minute |
||||
|
) |
||||
|
|
||||
|
// consoleOnlyRoutes 只允许 console 携带合法签名调用的路由(网关据此免登录放行 + 强制验签)。
|
||||
|
// 不放进 gateway.yaml 的 WhiteList——那是"免登录公开接口",语义完全不同。
|
||||
|
// 挂在 api 模块上:api 服务是业务侧的管理入口,且与其它业务服务同在 rpcx 集群,可扇出下发。
|
||||
|
var consoleOnlyRoutes = map[string]bool{ |
||||
|
"api_reloadmoduleconfig": true, |
||||
|
"api_resetmoduleconfig": true, |
||||
|
} |
||||
|
|
||||
|
// IsConsoleOnlyRoute 判断某路由是否为「仅 console 可调、必须验签」的运维接口。
|
||||
|
func IsConsoleOnlyRoute(msgName string) bool { return consoleOnlyRoutes[msgName] } |
||||
|
|
||||
|
// ModuleResetAll 传给「重置」接口的特殊模块名,表示整库还原为配置文件初始值(而非单个模块)。
|
||||
|
const ModuleResetAll = "*" |
||||
|
|
||||
|
// SignConsoleCall 生成调用签名。ts 为秒级 unix 时间戳字符串;arg 为接口参数(无参数传 "")。
|
||||
|
// 签名同时绑定 msgName 与 arg——否则拿到一个"重置 email"的签名就能改成"重置 wechatpay"。
|
||||
|
func SignConsoleCall(key, ts, msgName, arg string) string { |
||||
|
mac := hmac.New(sha256.New, []byte(key)) |
||||
|
mac.Write([]byte(ts + "." + msgName + "." + arg)) |
||||
|
return hex.EncodeToString(mac.Sum(nil)) |
||||
|
} |
||||
|
|
||||
|
// VerifyConsoleCall 校验签名与时间戳。key 取 ${FIELD_ENCRYPT_KEY}。
|
||||
|
// 用 hmac.Equal 做常数时间比较,避免按字节比较泄漏信息。
|
||||
|
func VerifyConsoleCall(key, ts, sign, msgName, arg string, now time.Time) error { |
||||
|
if key == "" { |
||||
|
return errors.New("未配置 ${FIELD_ENCRYPT_KEY},无法校验 console 调用签名") |
||||
|
} |
||||
|
if ts == "" || sign == "" { |
||||
|
return fmt.Errorf("缺少 %s / %s 请求头", ConsoleTsHeader, ConsoleSignHeader) |
||||
|
} |
||||
|
sec, err := strconv.ParseInt(ts, 10, 64) |
||||
|
if err != nil { |
||||
|
return errors.New("时间戳格式错误") |
||||
|
} |
||||
|
if d := now.Sub(time.Unix(sec, 0)); d > consoleSignSkew || d < -consoleSignSkew { |
||||
|
return errors.New("时间戳超出允许偏差(检查两端时钟)") |
||||
|
} |
||||
|
expect := SignConsoleCall(key, ts, msgName, arg) |
||||
|
if !hmac.Equal([]byte(sign), []byte(expect)) { |
||||
|
return errors.New("签名不匹配(检查 console 与本服务的 ${FIELD_ENCRYPT_KEY} 是否一致)") |
||||
|
} |
||||
|
return nil |
||||
|
} |
||||
@ -0,0 +1,76 @@ |
|||||
|
package comm |
||||
|
|
||||
|
import ( |
||||
|
"strconv" |
||||
|
"testing" |
||||
|
"time" |
||||
|
) |
||||
|
|
||||
|
func TestVerifyConsoleCall(t *testing.T) { |
||||
|
const key = "SzwdcXrtrmJMRZkZtQuzhAETOIJnFzMH" |
||||
|
const route = "user_login" |
||||
|
now := time.Unix(1_800_000_000, 0) |
||||
|
ts := strconv.FormatInt(now.Unix(), 10) |
||||
|
good := SignConsoleCall(key, ts, route, "") |
||||
|
// 翻转末位(末位本就可能是 '0',直接拼 "0" 会得到与原串相同的"篡改"签名,测试就成了假阳性)
|
||||
|
tampered := good[:len(good)-1] + map[bool]string{true: "1", false: "0"}[good[len(good)-1] == '0'] |
||||
|
|
||||
|
if err := VerifyConsoleCall(key, ts, good, route, "", now); err != nil { |
||||
|
t.Fatalf("合法签名应通过: %v", err) |
||||
|
} |
||||
|
// 时钟小幅偏差仍应通过(两机时钟不会完全一致)
|
||||
|
if err := VerifyConsoleCall(key, ts, good, route, "", now.Add(time.Minute)); err != nil { |
||||
|
t.Fatalf("1 分钟偏差应通过: %v", err) |
||||
|
} |
||||
|
|
||||
|
cases := []struct { |
||||
|
name string |
||||
|
key, ts, sign, rt, arg string |
||||
|
at time.Time |
||||
|
}{ |
||||
|
{"密钥不一致(两端 FIELD_ENCRYPT_KEY 不同)", "other-key", ts, good, route, "", now}, |
||||
|
{"签名被篡改", key, ts, tampered, route, "", now}, |
||||
|
{"签名绑定的是别的接口(防止换路由重放)", key, ts, SignConsoleCall(key, ts, "api_other", ""), route, "", now}, |
||||
|
{"时间戳过期", key, ts, good, route, "", now.Add(10 * time.Minute)}, |
||||
|
{"时间戳来自未来", key, ts, good, route, "", now.Add(-10 * time.Minute)}, |
||||
|
{"缺时间戳", key, "", good, route, "", now}, |
||||
|
{"缺签名", key, ts, "", route, "", now}, |
||||
|
{"未配置密钥", "", ts, good, route, "", now}, |
||||
|
{"时间戳非数字", key, "abc", good, route, "", now}, |
||||
|
} |
||||
|
for _, c := range cases { |
||||
|
t.Run(c.name, func(t *testing.T) { |
||||
|
if err := VerifyConsoleCall(c.key, c.ts, c.sign, c.rt, c.arg, c.at); err == nil { |
||||
|
t.Fatal("应当被拒绝") |
||||
|
} |
||||
|
}) |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
// 路由白名单与实际注册的接口必须对得上,否则网关要么放行了不该放行的,要么把 console 挡在门外。
|
||||
|
func TestConsoleOnlyRoutes(t *testing.T) { |
||||
|
if !IsConsoleOnlyRoute("api_reloadmoduleconfig") { |
||||
|
t.Fatal("重载接口必须是 console 专用路由(否则会裸奔在公网上)") |
||||
|
} |
||||
|
if !IsConsoleOnlyRoute("api_resetmoduleconfig") { |
||||
|
t.Fatal("重置接口必须是 console 专用路由") |
||||
|
} |
||||
|
if IsConsoleOnlyRoute("user_login") { |
||||
|
t.Fatal("普通用户接口不该被当作 console 专用路由") |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
// 参数(模块名)必须被签名绑定:否则截获一个"重置 email"的请求,改个头就能重置 wechatpay。
|
||||
|
func TestSignBindsArg(t *testing.T) { |
||||
|
const key, route = "k", "api_resetmoduleconfig" |
||||
|
now := time.Unix(1_800_000_000, 0) |
||||
|
ts := strconv.FormatInt(now.Unix(), 10) |
||||
|
signEmail := SignConsoleCall(key, ts, route, "email") |
||||
|
|
||||
|
if err := VerifyConsoleCall(key, ts, signEmail, route, "email", now); err != nil { |
||||
|
t.Fatalf("原参数应通过: %v", err) |
||||
|
} |
||||
|
if err := VerifyConsoleCall(key, ts, signEmail, route, "wechatpay", now); err == nil { |
||||
|
t.Fatal("参数被换成别的模块必须拒绝") |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,96 @@ |
|||||
|
package comm |
||||
|
|
||||
|
import ( |
||||
|
"strings" |
||||
|
"testing" |
||||
|
) |
||||
|
|
||||
|
// 指纹的唯一用途是「两端比对」:实际 AES 密钥相同 → 指纹必相同;不同 → 指纹必不同。
|
||||
|
func TestKeyFingerprint_同密钥同指纹_异密钥异指纹(t *testing.T) { |
||||
|
a := KeyFingerprint("SzwdcXrtrmJMRZkZtQuzhAETOIJnFzMH") |
||||
|
if a != KeyFingerprint("SzwdcXrtrmJMRZkZtQuzhAETOIJnFzMH") { |
||||
|
t.Fatal("同一把 key 两次指纹不一致") |
||||
|
} |
||||
|
if a == KeyFingerprint("") { |
||||
|
t.Fatal("真实 key 与空 key 指纹相同") |
||||
|
} |
||||
|
// 不得泄露密钥本身。
|
||||
|
if strings.Contains(a, "Szwdc") { |
||||
|
t.Fatalf("指纹泄露了密钥内容: %s", a) |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
// normAESKey 截断到 32 字节:前 32 字节相同的两把 key 实为同一把 AES 密钥,指纹须相同,
|
||||
|
// 否则两端会因"指纹不同"去排查一个根本不存在的问题。
|
||||
|
func TestKeyFingerprint_超长截断后等价(t *testing.T) { |
||||
|
k32 := "SzwdcXrtrmJMRZkZtQuzhAETOIJnFzMH" // 恰好 32 字节
|
||||
|
long := k32 + "这些字节根本不参与加密" |
||||
|
if KeyFingerprint(k32) == KeyFingerprint(long) { |
||||
|
// 指纹的哈希部分应相同(实际 AES key 相同),但备注不同,故比较哈希前缀。
|
||||
|
t.Fatal("期望备注不同") |
||||
|
} |
||||
|
h1, h2 := strings.Fields(KeyFingerprint(k32))[0], strings.Fields(KeyFingerprint(long))[0] |
||||
|
if h1 != h2 { |
||||
|
t.Fatalf("前 32 字节相同的 key 实际 AES 密钥相同,指纹哈希应一致: %s vs %s", h1, h2) |
||||
|
} |
||||
|
if !strings.Contains(KeyFingerprint(long), "被丢弃") { |
||||
|
t.Fatal("超长 key 应提示截断") |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
func TestKeyFingerprint_标注常见事故(t *testing.T) { |
||||
|
if !strings.Contains(KeyFingerprint(""), "未配置") { |
||||
|
t.Error("空 key 应提示未配置") |
||||
|
} |
||||
|
if !strings.Contains(KeyFingerprint("short-key"), "补零") { |
||||
|
t.Error("不足 32 字节应提示补零") |
||||
|
} |
||||
|
// .env 存成 CRLF 时,值尾部会多一个 \r,肉眼与正确值毫无差别。
|
||||
|
if !strings.Contains(KeyFingerprint("SzwdcXrtrmJMRZkZtQuzhAETOIJnFzM\r"), "CR") { |
||||
|
t.Error("尾部 CR 应被标注") |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
// 事故复现:FIELD_ENCRYPT_KEY 进 env 模板之前,业务服务用空串(→全零密钥)把配置 seed 进了库。
|
||||
|
// 后来 env 配上了真 key,解密失败——诊断必须指出"这是旧 key 加密的存量数据",而不是让人去查 env。
|
||||
|
func TestDiagnoseDecryptFailure_识别空key加密的存量数据(t *testing.T) { |
||||
|
const realKey = "SzwdcXrtrmJMRZkZtQuzhAETOIJnFzMH" |
||||
|
legacyCT, err := Encrypt("", "wx-appsecret-plaintext") // 旧进程用空 key 加密
|
||||
|
if err != nil { |
||||
|
t.Fatal(err) |
||||
|
} |
||||
|
if _, err = Decrypt(realKey, legacyCT); err == nil { |
||||
|
t.Fatal("前提不成立:真 key 不该能解开空 key 的密文") |
||||
|
} |
||||
|
|
||||
|
d := DiagnoseDecryptFailure(realKey, legacyCT) |
||||
|
if !strings.Contains(d, "全零密钥") { |
||||
|
t.Errorf("应指认出空 key,实际: %s", d) |
||||
|
} |
||||
|
if strings.Contains(d, "wx-appsecret-plaintext") { |
||||
|
t.Fatalf("诊断泄露了明文: %s", d) |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
// .env 存成 CRLF:值尾部混入 \r,两端"看起来"配的是同一把 key,实际差一个字节。
|
||||
|
func TestDiagnoseDecryptFailure_识别尾部CR(t *testing.T) { |
||||
|
const clean = "SzwdcXrtrmJMRZkZtQuzhAETOIJnFzMH" |
||||
|
ct, err := Encrypt(clean, "some-secret") |
||||
|
if err != nil { |
||||
|
t.Fatal(err) |
||||
|
} |
||||
|
d := DiagnoseDecryptFailure(clean+"\r", ct) // 本服务读到的 key 尾部多了 \r
|
||||
|
if !strings.Contains(d, "多余字符") { |
||||
|
t.Errorf("应指认出尾部混入字符,实际: %s", d) |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
func TestDiagnoseDecryptFailure_未知key不瞎猜(t *testing.T) { |
||||
|
ct, err := Encrypt("some-unknown-key-nobody-has-seen", "x") |
||||
|
if err != nil { |
||||
|
t.Fatal(err) |
||||
|
} |
||||
|
if d := DiagnoseDecryptFailure("SzwdcXrtrmJMRZkZtQuzhAETOIJnFzMH", ct); !strings.Contains(d, "未知的 key") { |
||||
|
t.Errorf("无线索时应如实说明,实际: %s", d) |
||||
|
} |
||||
|
} |
||||
@ -1 +1 @@ |
|||||
2026/07/09 15:12:59.319 ERROR comm/moduleconfig.go:200 [ModuleConfig] 密钥字段解密失败,已跳过该字段(保留默认值) module=wechatpay key=PrivateKeyPath err=GCM 解密失败: cipher: message authentication failed —— 检查本服务 ${FIELD_ENCRYPT_KEY} 是否与 console 完全一致 |
2026/07/10 10:56:05.452 ERROR comm/moduleconfig.go:202 [ModuleConfig] 密钥字段解密失败,已跳过该字段(保留默认值) module=wechatpay key=PrivateKeyPath err=GCM 解密失败: cipher: message authentication failed 本服务key指纹=sha256:19731e34cd8c [不足 32 字节,已右侧补零(原长 30)] 【诊断】密文是用「空 key(FIELD_ENCRYPT_KEY 未注入时补零出的全零密钥)」(指纹 sha256:17a63bfb3b5b [!!未配置 FIELD_ENCRYPT_KEY,正在使用全零密钥])加密的存量数据,与本服务当前 key 不是同一把 —— 到 console 后台把该字段重新填一遍保存,即可用当前 key 重新加密落库 |
||||
|
|||||
@ -0,0 +1,310 @@ |
|||||
|
package comm |
||||
|
|
||||
|
import ( |
||||
|
"errors" |
||||
|
"fmt" |
||||
|
"sync" |
||||
|
|
||||
|
"yunyan/lego/sys/log" |
||||
|
"yunyan/lego/sys/mysql" |
||||
|
) |
||||
|
|
||||
|
// 业务功能配置的「热重载」——console 后台改完配置广播事件,业务服务据此重读业务库、
|
||||
|
// 热替换对应的 sys 客户端,免去重启。
|
||||
|
//
|
||||
|
// 能不能热更取决于该 sys 持有什么:
|
||||
|
// - 凭据类客户端(email/sms/…):只是拿着 key/密码的 HTTP/SMTP 客户端,可整体替换 → applied
|
||||
|
// - 未注册热重载函数的模块:如实回 restart_required,绝不谎报成功
|
||||
|
// - 重载失败:保留旧客户端与旧配置继续服务 → failed
|
||||
|
//
|
||||
|
// 与 LoadOrSeedModuleConfig 的关键差别:本函数**只读覆盖、绝不 seed**——
|
||||
|
// 运行时不该因为某模块库里没配就把 yaml 默认值写进库。
|
||||
|
|
||||
|
var ( |
||||
|
reloaderMu sync.RWMutex |
||||
|
// reloaders SysKey → 热重载函数。由业务服务启动时 RegisterModuleReloader 注册。
|
||||
|
reloaders = map[string]func(map[string]interface{}) error{} |
||||
|
) |
||||
|
|
||||
|
// RegisterModuleReloader 注册某 SysKey(见 ModuleDef.SysKey) 的热重载函数。
|
||||
|
// 约定:实现必须「校验失败返回 error 且保留旧客户端」,绝不 panic——
|
||||
|
// 后台一次误保存不该打挂线上进程。
|
||||
|
func RegisterModuleReloader(sysKey string, fn func(map[string]interface{}) error) { |
||||
|
reloaderMu.Lock() |
||||
|
reloaders[sysKey] = fn |
||||
|
reloaderMu.Unlock() |
||||
|
} |
||||
|
|
||||
|
func lookupReloader(sysKey string) (func(map[string]interface{}) error, bool) { |
||||
|
reloaderMu.RLock() |
||||
|
defer reloaderMu.RUnlock() |
||||
|
fn, ok := reloaders[sysKey] |
||||
|
return fn, ok |
||||
|
} |
||||
|
|
||||
|
// HasModuleReloaders 本服务是否持有任何可热重载的业务凭据客户端。
|
||||
|
// gateway/timer/mcp 一个都不注册——它们收到下发时如实回「本服务不承载业务配置」,而非报错。
|
||||
|
func HasModuleReloaders() bool { |
||||
|
reloaderMu.RLock() |
||||
|
defer reloaderMu.RUnlock() |
||||
|
return len(reloaders) > 0 |
||||
|
} |
||||
|
|
||||
|
// ───────── 集群内下发的 rpc 载荷(rpcx 用 MsgPack 序列化,普通 struct 即可,无需 pb) ─────────
|
||||
|
|
||||
|
// ReloadModuleConfigReq api 服务 → 各业务服务:请重载业务功能配置。
|
||||
|
type ReloadModuleConfigReq struct { |
||||
|
Reason string // 触发原因(save / reset),仅用于日志
|
||||
|
} |
||||
|
|
||||
|
// ResetModuleConfigReq api 服务 → 配置属主(home):把某模块还原为你 yaml 里的初始值。
|
||||
|
type ResetModuleConfigReq struct { |
||||
|
Module string |
||||
|
} |
||||
|
|
||||
|
// ModuleConfigResp 单个服务对下发的应答。Err 非空表示该服务处理失败——
|
||||
|
// 用字段而非 rpc error 承载,好让 api 汇总时能区分「某个服务失败」与「整体调用失败」。
|
||||
|
type ModuleConfigResp struct { |
||||
|
Service string // 服务名(home/api…)
|
||||
|
Instance string // <服务名>@<主机名>
|
||||
|
Skipped bool // 本服务不承载业务配置,什么也没做
|
||||
|
Results []ModuleApplyResult // 逐模块三态
|
||||
|
Err string |
||||
|
} |
||||
|
|
||||
|
// reloadContext 执行一次重载所需的全部依赖,由业务服务启动时登记。
|
||||
|
// 有了它,HTTP 重载接口(处在某个业务模块里)不必再层层拿到服务的 Settings 与库句柄。
|
||||
|
type reloadContext struct { |
||||
|
db mysql.ISys |
||||
|
encKey string |
||||
|
sys map[string]map[string]interface{} |
||||
|
} |
||||
|
|
||||
|
var ( |
||||
|
reloadCtxMu sync.RWMutex |
||||
|
reloadCtx *reloadContext |
||||
|
) |
||||
|
|
||||
|
// SetModuleReloadContext 登记重载上下文(业务服务启动、各 sys OnInit 之后调用)。
|
||||
|
func SetModuleReloadContext(db mysql.ISys, encKey string, sys map[string]map[string]interface{}) { |
||||
|
reloadCtxMu.Lock() |
||||
|
reloadCtx = &reloadContext{db: db, encKey: encKey, sys: sys} |
||||
|
reloadCtxMu.Unlock() |
||||
|
} |
||||
|
|
||||
|
// DoReloadModuleConfig 用已登记的上下文执行一次业务配置热重载。
|
||||
|
// 未登记(本服务不承载业务配置)时返回 error,由调用方如实上报,绝不假装成功。
|
||||
|
func DoReloadModuleConfig() ([]ModuleApplyResult, error) { |
||||
|
reloadCtxMu.RLock() |
||||
|
ctx := reloadCtx |
||||
|
reloadCtxMu.RUnlock() |
||||
|
if ctx == nil { |
||||
|
return nil, errors.New("本服务未登记业务配置重载上下文") |
||||
|
} |
||||
|
return ReloadModuleConfig(ctx.db, ctx.encKey, ctx.sys), nil |
||||
|
} |
||||
|
|
||||
|
// ───────────────────────── 重置为「配置文件初始值」 ─────────────────────────
|
||||
|
//
|
||||
|
// 业务库里的配置被改脏、又没有备份时的兜底:把某模块整体还原成本服务 confs/*.yaml 里的初始值。
|
||||
|
//
|
||||
|
// 数据来源只能是「启动时、库值覆盖之前」抓的 yaml 快照——一旦 LoadOrSeedModuleConfig 跑完,
|
||||
|
// 内存里的 Sys 就已经是库值了,那时再读就是把脏数据又写回去。
|
||||
|
|
||||
|
var ( |
||||
|
fileCfgMu sync.RWMutex |
||||
|
fileCfg map[string]map[string]interface{} // SysKey → yaml 原值副本
|
||||
|
) |
||||
|
|
||||
|
// snapshotFileModuleConfig 深拷贝各模块的 yaml 原值。由 LoadOrSeedModuleConfig 在覆盖 Sys 之前调用。
|
||||
|
func snapshotFileModuleConfig(sys map[string]map[string]interface{}) { |
||||
|
snap := make(map[string]map[string]interface{}, len(ModuleCatalog)) |
||||
|
for i := range ModuleCatalog { |
||||
|
sec := sys[ModuleCatalog[i].SysKey] |
||||
|
if sec == nil { |
||||
|
continue |
||||
|
} |
||||
|
cp := make(map[string]interface{}, len(sec)) |
||||
|
for k, v := range sec { |
||||
|
cp[k] = v |
||||
|
} |
||||
|
snap[ModuleCatalog[i].SysKey] = cp |
||||
|
} |
||||
|
fileCfgMu.Lock() |
||||
|
fileCfg = snap |
||||
|
fileCfgMu.Unlock() |
||||
|
} |
||||
|
|
||||
|
func fileSnapshot(sysKey string) map[string]interface{} { |
||||
|
fileCfgMu.RLock() |
||||
|
defer fileCfgMu.RUnlock() |
||||
|
return fileCfg[sysKey] |
||||
|
} |
||||
|
|
||||
|
// DoResetModuleConfigToFile 把某模块在业务库里的配置整体还原为本服务配置文件(yaml)的初始值。
|
||||
|
// yaml 里没有该模块(或本服务未登记上下文)时**不做任何改动**并返回 error——
|
||||
|
// 否则会把该模块清空后无值可填,比"脏数据"更糟。
|
||||
|
func DoResetModuleConfigToFile(module string) ([]ModuleApplyResult, error) { |
||||
|
ctx := currentReloadCtx() |
||||
|
if ctx == nil { |
||||
|
return nil, errors.New("本服务未登记业务配置重载上下文") |
||||
|
} |
||||
|
def := FindModuleDef(module) |
||||
|
if def == nil { |
||||
|
return nil, fmt.Errorf("未知模块: %s", module) |
||||
|
} |
||||
|
if len(fileSnapshot(def.SysKey)) == 0 { |
||||
|
return nil, fmt.Errorf("本服务的配置文件里没有 %s(%s) 的配置,无法还原", def.Name, def.Module) |
||||
|
} |
||||
|
res, err := resetOneModuleToFile(ctx, def) |
||||
|
if err != nil { |
||||
|
return nil, err |
||||
|
} |
||||
|
return []ModuleApplyResult{res}, nil |
||||
|
} |
||||
|
|
||||
|
// DoResetAllModuleConfigToFile 把业务库里**所有**模块的配置整体还原为配置文件(yaml)初始值——
|
||||
|
// 用于升级/迁移把库写脏后的整体恢复。只处理 yaml 里确有配置的模块(快照非空);
|
||||
|
// yaml 未定义的模块保持不动(避免清空后无值可填)。
|
||||
|
//
|
||||
|
// 单个模块失败不中断其它模块:把失败作为该模块的 result 收集,尽最大努力恢复整库。
|
||||
|
// ⚠️ 破坏性:会覆盖后台手工保存过的值(如把 email 打回 yaml 的默认发信配置),调用方须二次确认。
|
||||
|
func DoResetAllModuleConfigToFile() ([]ModuleApplyResult, error) { |
||||
|
ctx := currentReloadCtx() |
||||
|
if ctx == nil { |
||||
|
return nil, errors.New("本服务未登记业务配置重载上下文") |
||||
|
} |
||||
|
results := make([]ModuleApplyResult, 0, len(ModuleCatalog)) |
||||
|
for i := range ModuleCatalog { |
||||
|
def := &ModuleCatalog[i] |
||||
|
if len(fileSnapshot(def.SysKey)) == 0 { |
||||
|
continue // yaml 没定义这个模块 → 不动它
|
||||
|
} |
||||
|
res, err := resetOneModuleToFile(ctx, def) |
||||
|
if err != nil { |
||||
|
log.Errorf("[ModuleReset] 全量重置中 %s 失败: %v", def.Module, err) |
||||
|
res = ModuleApplyResult{Module: def.Module, Status: ApplyStatusFailed, Err: err.Error()} |
||||
|
} |
||||
|
results = append(results, res) |
||||
|
} |
||||
|
if len(results) == 0 { |
||||
|
return nil, errors.New("配置文件里没有任何业务模块配置,无法还原") |
||||
|
} |
||||
|
log.Infof("[ModuleReset] 已全量还原为配置文件初始值 modules=%d", len(results)) |
||||
|
return results, nil |
||||
|
} |
||||
|
|
||||
|
func currentReloadCtx() *reloadContext { |
||||
|
reloadCtxMu.RLock() |
||||
|
defer reloadCtxMu.RUnlock() |
||||
|
return reloadCtx |
||||
|
} |
||||
|
|
||||
|
// resetOneModuleToFile 删该模块全部行 → 按 yaml 快照重写(密钥加密、带时间戳) → 覆盖 Sys → 触发热重载。
|
||||
|
// 调用方已保证 fileSnapshot(def.SysKey) 非空。
|
||||
|
func resetOneModuleToFile(ctx *reloadContext, def *ModuleDef) (ModuleApplyResult, error) { |
||||
|
snap := fileSnapshot(def.SysKey) |
||||
|
if err := ctx.db.Delete(TableAppModuleConfig, "module=?", def.Module); err != nil { |
||||
|
return ModuleApplyResult{}, fmt.Errorf("清除 %s 旧配置失败: %w", def.Module, err) |
||||
|
} |
||||
|
if err := seedModuleRows(ctx.db, ctx.encKey, def, snap); err != nil { |
||||
|
return ModuleApplyResult{}, fmt.Errorf("写入 %s 配置文件初始值失败(该模块配置现为空,请重试): %w", def.Module, err) |
||||
|
} |
||||
|
|
||||
|
// Sys 与库保持一致,再让客户端切到初始值。
|
||||
|
next := make(map[string]interface{}, len(snap)) |
||||
|
for k, v := range snap { |
||||
|
next[k] = v |
||||
|
} |
||||
|
ctx.sys[def.SysKey] = next |
||||
|
|
||||
|
res := ModuleApplyResult{Module: def.Module, Status: ApplyStatusRestartRequired} |
||||
|
if fn, ok := lookupReloader(def.SysKey); ok { |
||||
|
if err := fn(next); err != nil { |
||||
|
log.Errorf("[ModuleReset] 还原后热重载失败 module=%s err=%v", def.Module, err) |
||||
|
res = ModuleApplyResult{Module: def.Module, Status: ApplyStatusFailed, Err: err.Error()} |
||||
|
} else { |
||||
|
res = ModuleApplyResult{Module: def.Module, Status: ApplyStatusApplied} |
||||
|
} |
||||
|
} |
||||
|
log.Infof("[ModuleReset] 已还原为配置文件初始值 module=%s status=%s", def.Module, res.Status) |
||||
|
return res, nil |
||||
|
} |
||||
|
|
||||
|
// ReloadModuleConfig 重读业务库 app_module_config,对每个「库里有配置」的模块尝试热重载。
|
||||
|
//
|
||||
|
// db 应用业务库(与 LoadOrSeedModuleConfig 同一个)
|
||||
|
// encKey ${FIELD_ENCRYPT_KEY},须与 console 一致
|
||||
|
// sys GetSettings().Sys(map 引用);仅在该模块重载成功/待重启时才提交新值
|
||||
|
//
|
||||
|
// 返回每个模块的结果,供上层组装 ConfigApplyAck 回给 console。
|
||||
|
// 单个模块失败不影响其它模块,也不会改动它自己原有的客户端。
|
||||
|
func ReloadModuleConfig(db mysql.ISys, encKey string, sys map[string]map[string]interface{}) []ModuleApplyResult { |
||||
|
results := make([]ModuleApplyResult, 0, len(ModuleCatalog)) |
||||
|
if db == nil { |
||||
|
return append(results, ModuleApplyResult{Module: "*", Status: ApplyStatusFailed, Err: "业务库未连接"}) |
||||
|
} |
||||
|
rows := make([]*AppModuleConfig, 0) |
||||
|
if err := db.Find(TableAppModuleConfig, &rows, ""); err != nil { |
||||
|
return append(results, ModuleApplyResult{Module: "*", Status: ApplyStatusFailed, Err: "读取业务配置失败: " + err.Error()}) |
||||
|
} |
||||
|
byModule := make(map[string][]*AppModuleConfig) |
||||
|
for _, r := range rows { |
||||
|
byModule[r.Module] = append(byModule[r.Module], r) |
||||
|
} |
||||
|
|
||||
|
for i := range ModuleCatalog { |
||||
|
def := &ModuleCatalog[i] |
||||
|
existing := byModule[def.Module] |
||||
|
if len(existing) == 0 { |
||||
|
continue // 库里没这个模块的配置 → 不动它(也不 seed)
|
||||
|
} |
||||
|
|
||||
|
// 先在副本上组装新配置:重载失败时 sys 保持原样,避免留下一份「生效不了的」脏配置。
|
||||
|
cur := sys[def.SysKey] |
||||
|
next := make(map[string]interface{}, len(cur)+len(existing)) |
||||
|
for k, v := range cur { |
||||
|
next[k] = v |
||||
|
} |
||||
|
// **fail-closed,且以模块为原子**:只要本模块有任一密钥字段解不开,就整模块放弃重载。
|
||||
|
// 半新半旧的凭据组合比原地不动更危险;而把解不开的密文当明文塞进 next,等于让
|
||||
|
// email.Reload 拿密文当 SMTP 密码、让 wechatpay 拿密文当私钥路径去 open——
|
||||
|
// 后者正是 home 进程被 log.Fatal 杀死、5-in-1 容器 crash loop、全站 502 的那条路。
|
||||
|
decryptFailed := 0 |
||||
|
for _, r := range existing { |
||||
|
val := r.Value |
||||
|
if r.Encrypted && val != "" { |
||||
|
plain, e := Decrypt(encKey, val) |
||||
|
if e != nil { |
||||
|
decryptFailed++ |
||||
|
log.Errorf("[ModuleReload] 密钥字段解密失败 module=%s key=%s err=%v 本服务key指纹=%s %s", |
||||
|
r.Module, r.Key, e, KeyFingerprint(encKey), DiagnoseDecryptFailure(encKey, r.Value)) |
||||
|
continue |
||||
|
} |
||||
|
val = plain |
||||
|
} |
||||
|
next[r.Key] = coerceVal(val, cur[r.Key]) |
||||
|
} |
||||
|
if decryptFailed > 0 { |
||||
|
results = append(results, ModuleApplyResult{Module: def.Module, Status: ApplyStatusFailed, |
||||
|
Err: fmt.Sprintf("%d 个密钥字段解密失败,已跳过本模块重载(继续使用原有配置)——见本服务日志里的 key 指纹与诊断", decryptFailed)}) |
||||
|
continue |
||||
|
} |
||||
|
|
||||
|
fn, ok := lookupReloader(def.SysKey) |
||||
|
if !ok { |
||||
|
sys[def.SysKey] = next // 下次重启即生效
|
||||
|
results = append(results, ModuleApplyResult{Module: def.Module, Status: ApplyStatusRestartRequired}) |
||||
|
continue |
||||
|
} |
||||
|
if err := fn(next); err != nil { |
||||
|
log.Errorf("[ModuleReload] 模块热重载失败 module=%s err=%v(已保留旧配置继续服务)", def.Module, err) |
||||
|
results = append(results, ModuleApplyResult{Module: def.Module, Status: ApplyStatusFailed, Err: err.Error()}) |
||||
|
continue |
||||
|
} |
||||
|
sys[def.SysKey] = next |
||||
|
log.Infof("[ModuleReload] 模块热重载成功 module=%s", def.Module) |
||||
|
results = append(results, ModuleApplyResult{Module: def.Module, Status: ApplyStatusApplied}) |
||||
|
} |
||||
|
return results |
||||
|
} |
||||
Binary file not shown.
@ -0,0 +1,107 @@ |
|||||
|
package api |
||||
|
|
||||
|
import ( |
||||
|
"context" |
||||
|
"encoding/json" |
||||
|
"os" |
||||
|
"time" |
||||
|
|
||||
|
"yunyan/comm" |
||||
|
"yunyan/lego/sys/log" |
||||
|
"yunyan/pb" |
||||
|
) |
||||
|
|
||||
|
// 业务功能配置的「管理入口」——console 只跟本服务打交道,由本服务在 ETCD/rpcx 集群里扇出。
|
||||
|
//
|
||||
|
// 为什么是 api 服务:它是业务侧对外的管理入口,且与其它业务服务同在一个 rpcx 集群,
|
||||
|
// 能直接把配置变更下发给每一个持有业务凭据的服务(home 有 email/sms/auth/pay,api 自己也有 email)。
|
||||
|
// 只重载其中一个,另一个照旧用着旧凭据——这正是之前把重载塞进 home 的 user 模块所犯的错。
|
||||
|
//
|
||||
|
// 鉴权:不走用户登录,由网关按 comm.IsConsoleOnlyRoute 强制校验 ${FIELD_ENCRYPT_KEY} 的 HMAC 签名。
|
||||
|
// 返回 []byte 原样 JSON(反射路由允许),故无需为本接口新增 pb 消息类型。
|
||||
|
|
||||
|
const fanoutTimeout = 8 * time.Second |
||||
|
|
||||
|
// ReloadModuleConfig 让本部署内所有业务服务重读业务库、热替换各自的 sys 客户端。
|
||||
|
// 路由 api_reloadmoduleconfig → POST /api/api/api_reloadmoduleconfig
|
||||
|
func (this *apiComp) ReloadModuleConfig(session comm.IUserSession, req *pb.Rpc_EmptyReq) ([]byte, *pb.ErrorData) { |
||||
|
return this.marshalReport(this.fanoutReload("save")) |
||||
|
} |
||||
|
|
||||
|
// ResetModuleConfig 把某模块还原为**配置属主(home)** 配置文件里的初始值,随后让全体重读库。
|
||||
|
// 目标模块经 X-Console-Arg 请求头传入(已并入 HMAC 签名,网关校验后搬进 Meta)。
|
||||
|
// 路由 api_resetmoduleconfig → POST /api/api/api_resetmoduleconfig
|
||||
|
func (this *apiComp) ResetModuleConfig(session comm.IUserSession, req *pb.Rpc_EmptyReq) ([]byte, *pb.ErrorData) { |
||||
|
module := session.GetMateToString(comm.SessionMeta_ConsoleArg) |
||||
|
if module == "" { |
||||
|
return nil, &pb.ErrorData{Code: pb.ErrorCode_ReqParameterError, Message: "缺少要重置的模块名"} |
||||
|
} |
||||
|
// "*" = 整库还原;否则必须是目录里已知的模块。
|
||||
|
if module != comm.ModuleResetAll && comm.FindModuleDef(module) == nil { |
||||
|
return nil, &pb.ErrorData{Code: pb.ErrorCode_ReqParameterError, Message: "未知模块: " + module} |
||||
|
} |
||||
|
|
||||
|
// ① 只让属主用它自己的 yaml 重写库。失败就直接返回——库没动过,不该再去广播重载。
|
||||
|
ctx, cancel := context.WithTimeout(context.Background(), fanoutTimeout) |
||||
|
defer cancel() |
||||
|
var rr comm.ModuleConfigResp |
||||
|
if err := this.service.RpcCall(ctx, comm.ModuleConfigOwner, |
||||
|
string(comm.Rpc_ResetModuleConfig), &comm.ResetModuleConfigReq{Module: module}, &rr); err != nil { |
||||
|
return nil, &pb.ErrorData{Code: pb.ErrorCode_SystemError, Message: "调用配置属主(" + comm.ModuleConfigOwner + ")失败: " + err.Error()} |
||||
|
} |
||||
|
if rr.Err != "" { |
||||
|
return nil, &pb.ErrorData{Code: pb.ErrorCode_SystemError, Message: "还原失败: " + rr.Err} |
||||
|
} |
||||
|
|
||||
|
// ② 库已是初始值,通知全体重读并热换客户端。
|
||||
|
return this.marshalReport(this.fanoutReload("reset")) |
||||
|
} |
||||
|
|
||||
|
// fanoutReload 向每个持有业务凭据的服务下发重载,收集各自回执。
|
||||
|
//
|
||||
|
// 用 RpcCall 逐服务下发而非 RpcBroadcast:Broadcast 打得到某服务的全部实例,却拿不回逐实例的结构化结果,
|
||||
|
// 而"配置到底生效没有"正是这个功能的全部意义。代价是单服务多实例时只命中其中一个——
|
||||
|
// 当前每个部署每服务各一个实例;将来扩副本需改用 Broadcast 并接受更粗的回执。
|
||||
|
func (this *apiComp) fanoutReload(reason string) []comm.ConfigApplyAck { |
||||
|
ctx, cancel := context.WithTimeout(context.Background(), fanoutTimeout) |
||||
|
defer cancel() |
||||
|
|
||||
|
acks := make([]comm.ConfigApplyAck, 0, len(comm.ModuleConfigServices)) |
||||
|
now := time.Now().Unix() |
||||
|
for _, svc := range comm.ModuleConfigServices { |
||||
|
var resp comm.ModuleConfigResp |
||||
|
if err := this.service.RpcCall(ctx, svc, |
||||
|
string(comm.Rpc_ReloadModuleConfig), &comm.ReloadModuleConfigReq{Reason: reason}, &resp); err != nil { |
||||
|
// 调不通也要如实上报:这个服务仍在用旧配置,不能当作成功。
|
||||
|
log.Error("[ApiModuleConfig] 下发重载失败", log.Field{Key: "service", Value: svc}, log.Field{Key: "err", Value: err.Error()}) |
||||
|
acks = append(acks, comm.ConfigApplyAck{ |
||||
|
Instance: svc + "@unreachable", Kind: comm.ConfigKindModuleConfig, Ts: now, |
||||
|
Results: []comm.ModuleApplyResult{{Module: "*", Status: comm.ApplyStatusFailed, Err: "服务不可达: " + err.Error()}}, |
||||
|
}) |
||||
|
continue |
||||
|
} |
||||
|
if resp.Skipped { |
||||
|
continue // 该服务不持有业务凭据,没什么可重载的,不必出现在回执里
|
||||
|
} |
||||
|
ack := comm.ConfigApplyAck{Instance: resp.Instance, Kind: comm.ConfigKindModuleConfig, Results: resp.Results, Ts: now} |
||||
|
if resp.Err != "" { |
||||
|
ack.Results = append(ack.Results, comm.ModuleApplyResult{Module: "*", Status: comm.ApplyStatusFailed, Err: resp.Err}) |
||||
|
} |
||||
|
acks = append(acks, ack) |
||||
|
} |
||||
|
return acks |
||||
|
} |
||||
|
|
||||
|
func (this *apiComp) marshalReport(acks []comm.ConfigApplyAck) ([]byte, *pb.ErrorData) { |
||||
|
report := comm.ConfigApplyReport{ |
||||
|
AppName: os.Getenv("ANALYZE_APP_NAME"), |
||||
|
Kind: comm.ConfigKindModuleConfig, |
||||
|
Acks: acks, |
||||
|
Ts: time.Now().Unix(), |
||||
|
} |
||||
|
body, err := json.Marshal(&report) |
||||
|
if err != nil { |
||||
|
return nil, &pb.ErrorData{Code: pb.ErrorCode_SystemError, Message: "回执序列化失败: " + err.Error()} |
||||
|
} |
||||
|
return body, nil |
||||
|
} |
||||
@ -0,0 +1,94 @@ |
|||||
|
package console |
||||
|
|
||||
|
import ( |
||||
|
"bytes" |
||||
|
"encoding/json" |
||||
|
"fmt" |
||||
|
"io" |
||||
|
"net/http" |
||||
|
"strconv" |
||||
|
"strings" |
||||
|
"time" |
||||
|
|
||||
|
"yunyan/comm" |
||||
|
) |
||||
|
|
||||
|
// console → 业务部署的「热重载」调用。
|
||||
|
//
|
||||
|
// console 是单例服务(lego/base/single)、不接 ETCD/rpcx 集群,无法直接 RPC 业务服务,
|
||||
|
// 只能走该部署的公网网关入口(app_registry.base_url)。鉴权用 ${FIELD_ENCRYPT_KEY} 的 HMAC 签名,
|
||||
|
// 由网关校验(见 comm/consolesign.go)——复用这把两边本就必须一致的密钥,不再多引入一把要同步的。
|
||||
|
|
||||
|
const ( |
||||
|
reloadRoute = "api_reloadmoduleconfig" |
||||
|
resetRoute = "api_resetmoduleconfig" |
||||
|
reloadTimeout = 5 * time.Second |
||||
|
) |
||||
|
|
||||
|
var reloadClient = &http.Client{Timeout: reloadTimeout} |
||||
|
|
||||
|
// callModuleConfigReload 通知目标部署重读业务库并热替换客户端。
|
||||
|
func callModuleConfigReload(baseUrl, encKey string) (*comm.ConfigApplyReport, error) { |
||||
|
return callConsoleRoute(baseUrl, encKey, reloadRoute, "") |
||||
|
} |
||||
|
|
||||
|
// callModuleConfigReset 让目标部署把某模块的配置还原为它自己 confs/*.yaml 的初始值。
|
||||
|
// yaml 原值只存在于业务服务内部,故必须由它自己完成——明文密钥一步都不离开该服务。
|
||||
|
func callModuleConfigReset(baseUrl, encKey, module string) (*comm.ConfigApplyReport, error) { |
||||
|
return callConsoleRoute(baseUrl, encKey, resetRoute, module) |
||||
|
} |
||||
|
|
||||
|
// callConsoleRoute 调目标部署的 console 专用运维接口并返回它的回执。
|
||||
|
// base_url 未填 / 网络不通 / 服务未起 / 验签失败 → 返回 error;调用方须如实呈现失败,
|
||||
|
// 绝不能默认当成功。arg 会并入签名,网关校验后经 Meta 传给业务侧。
|
||||
|
func callConsoleRoute(baseUrl, encKey, route, arg string) (*comm.ConfigApplyReport, error) { |
||||
|
base := strings.TrimRight(strings.TrimSpace(baseUrl), "/") |
||||
|
if base == "" { |
||||
|
return nil, fmt.Errorf("该部署未填写对外地址(base_url),无法通知它") |
||||
|
} |
||||
|
if !strings.HasPrefix(base, "http://") && !strings.HasPrefix(base, "https://") { |
||||
|
base = "https://" + base |
||||
|
} |
||||
|
|
||||
|
ts := strconv.FormatInt(time.Now().Unix(), 10) |
||||
|
req, err := http.NewRequest(http.MethodPost, base+"/api/api/"+route, bytes.NewReader([]byte("{}"))) |
||||
|
if err != nil { |
||||
|
return nil, err |
||||
|
} |
||||
|
req.Header.Set("Content-Type", "application/json") |
||||
|
req.Header.Set(comm.ConsoleTsHeader, ts) |
||||
|
req.Header.Set(comm.ConsoleArgHeader, arg) |
||||
|
req.Header.Set(comm.ConsoleSignHeader, comm.SignConsoleCall(encKey, ts, route, arg)) |
||||
|
|
||||
|
resp, err := reloadClient.Do(req) |
||||
|
if err != nil { |
||||
|
return nil, err |
||||
|
} |
||||
|
defer resp.Body.Close() |
||||
|
body, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20)) |
||||
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 { |
||||
|
return nil, fmt.Errorf("重载接口 HTTP %d: %s", resp.StatusCode, truncate(string(body), 200)) |
||||
|
} |
||||
|
|
||||
|
// 成功时 api 服务返回 ConfigApplyReport 的原样 JSON;出错时网关/业务返回 HttpResult{code,msg}。
|
||||
|
// 用 Kind 是否落位来区分两者——HttpResult 里没有这个字段。
|
||||
|
var report comm.ConfigApplyReport |
||||
|
if e := json.Unmarshal(body, &report); e == nil && report.Kind != "" { |
||||
|
return &report, nil |
||||
|
} |
||||
|
var hr struct { |
||||
|
Code int `json:"code"` |
||||
|
Message string `json:"msg"` |
||||
|
} |
||||
|
if json.Unmarshal(body, &hr) == nil && hr.Message != "" { |
||||
|
return nil, fmt.Errorf("重载被拒绝(code=%d): %s", hr.Code, hr.Message) |
||||
|
} |
||||
|
return nil, fmt.Errorf("重载接口返回无法识别的响应: %s", truncate(string(body), 200)) |
||||
|
} |
||||
|
|
||||
|
func truncate(s string, n int) string { |
||||
|
if len(s) <= n { |
||||
|
return s |
||||
|
} |
||||
|
return s[:n] + "…" |
||||
|
} |
||||
@ -0,0 +1,101 @@ |
|||||
|
package svccfg |
||||
|
|
||||
|
import ( |
||||
|
"context" |
||||
|
"os" |
||||
|
|
||||
|
"yunyan/comm" |
||||
|
"yunyan/lego/core" |
||||
|
"yunyan/lego/core/cbase" |
||||
|
"yunyan/lego/sys/log" |
||||
|
) |
||||
|
|
||||
|
// configComp 注册集群内的配置下发处理器。args/reply 用普通 struct——rpcx 走 MsgPack,无需 pb 类型。
|
||||
|
type configComp struct { |
||||
|
cbase.ModuleCompBase |
||||
|
module *SvcCfg |
||||
|
service comm.IService |
||||
|
} |
||||
|
|
||||
|
func (this *configComp) Init(service core.IService, module core.IModule, comp core.IModuleComp, opt core.IModuleOptions) (err error) { |
||||
|
if err = this.ModuleCompBase.Init(service, module, comp, opt); err != nil { |
||||
|
return |
||||
|
} |
||||
|
this.module = module.(*SvcCfg) |
||||
|
this.service = service.(comm.IService) |
||||
|
return |
||||
|
} |
||||
|
|
||||
|
func (this *configComp) Start() (err error) { |
||||
|
if err = this.ModuleCompBase.Start(); err != nil { |
||||
|
return |
||||
|
} |
||||
|
if err = this.service.Register(string(comm.Rpc_ReloadModuleConfig), this.Rpc_ReloadModuleConfig); err != nil { |
||||
|
return |
||||
|
} |
||||
|
// reset 处理器各服务都注册,但只有属主会真的执行(见下)——注册在非属主上是为了给出明确的拒绝理由,
|
||||
|
// 而不是让调用方收到"方法不存在"这种含糊的 rpc 错误。
|
||||
|
if err = this.service.Register(string(comm.Rpc_ResetModuleConfig), this.Rpc_ResetModuleConfig); err != nil { |
||||
|
return |
||||
|
} |
||||
|
log.Infof("[SvcCfg] 业务配置下发处理器已注册 service=%s owner=%v reloaders=%v", |
||||
|
this.service.GetType(), this.service.GetType() == comm.ModuleConfigOwner, comm.HasModuleReloaders()) |
||||
|
return |
||||
|
} |
||||
|
|
||||
|
func (this *configComp) instance() string { |
||||
|
host, _ := os.Hostname() |
||||
|
if host == "" { |
||||
|
host = "unknown" |
||||
|
} |
||||
|
return this.service.GetType() + "@" + host |
||||
|
} |
||||
|
|
||||
|
// Rpc_ReloadModuleConfig 重读业务库、热替换本服务持有的 sys 客户端。
|
||||
|
// 失败信息放在 reply.Err 而非 rpc error:这样 api 汇总时能区分「某个服务重载失败」与「整条 rpc 调不通」。
|
||||
|
func (this *configComp) Rpc_ReloadModuleConfig(ctx context.Context, args *comm.ReloadModuleConfigReq, reply *comm.ModuleConfigResp) (err error) { |
||||
|
reply.Service, reply.Instance = this.service.GetType(), this.instance() |
||||
|
if !comm.HasModuleReloaders() { |
||||
|
reply.Skipped = true // 本服务不持有任何业务凭据客户端(如 gateway/timer),什么也不用做
|
||||
|
return |
||||
|
} |
||||
|
results, e := comm.DoReloadModuleConfig() |
||||
|
if e != nil { |
||||
|
reply.Err = e.Error() |
||||
|
return |
||||
|
} |
||||
|
reply.Results = results |
||||
|
log.Infof("[SvcCfg] 已重载业务配置 service=%s reason=%s modules=%d", reply.Service, args.Reason, len(results)) |
||||
|
return |
||||
|
} |
||||
|
|
||||
|
// Rpc_ResetModuleConfig 把某模块还原为**本服务配置文件**里的初始值。
|
||||
|
//
|
||||
|
// 只有配置属主(comm.ModuleConfigOwner=home)会执行:yaml 每服务一份,非属主的 yaml 不是权威初始值,
|
||||
|
// 让它们也写库会互相覆盖。还原后仍需由 api 扇出 reload,让其余服务重读库。
|
||||
|
func (this *configComp) Rpc_ResetModuleConfig(ctx context.Context, args *comm.ResetModuleConfigReq, reply *comm.ModuleConfigResp) (err error) { |
||||
|
reply.Service, reply.Instance = this.service.GetType(), this.instance() |
||||
|
if reply.Service != comm.ModuleConfigOwner { |
||||
|
reply.Err = "本服务(" + reply.Service + ")不是业务配置属主(" + comm.ModuleConfigOwner + "),不能从配置文件还原" |
||||
|
return |
||||
|
} |
||||
|
if args.Module == "" { |
||||
|
reply.Err = "缺少要还原的模块名" |
||||
|
return |
||||
|
} |
||||
|
// "*" = 整库还原(升级/迁移把库写脏后的整体恢复);否则只还原单个模块。
|
||||
|
var results []comm.ModuleApplyResult |
||||
|
var e error |
||||
|
if args.Module == comm.ModuleResetAll { |
||||
|
results, e = comm.DoResetAllModuleConfigToFile() |
||||
|
} else { |
||||
|
results, e = comm.DoResetModuleConfigToFile(args.Module) |
||||
|
} |
||||
|
if e != nil { |
||||
|
reply.Err = e.Error() |
||||
|
return |
||||
|
} |
||||
|
reply.Results = results |
||||
|
log.Infof("[SvcCfg] 已还原为配置文件初始值 service=%s module=%s", reply.Service, args.Module) |
||||
|
return |
||||
|
} |
||||
@ -0,0 +1,42 @@ |
|||||
|
/* |
||||
|
服务配置管理模块 —— 装在每个「持有业务凭据 sys 客户端」的业务服务里(当前 home / api)。 |
||||
|
|
||||
|
职责只有一个:接收集群内下发的「业务功能配置变更」通知,重载**本服务**持有的 sys 客户端。 |
||||
|
|
||||
|
为什么独立成模块而不是塞进 user/api: |
||||
|
- 它是运维/管理面,不该寄生在终端用户接口模块里(权限与语义边界完全不同); |
||||
|
- 每个业务服务都持有各自的 sys 客户端(home 有 email/sms/auth/pay,api 也有 email), |
||||
|
只重载其中一个,另一个照旧用着旧凭据——必须每个服务各自重载。 |
||||
|
|
||||
|
下发路径:console --HTTPS+HMAC--> gateway --rpcx--> api 服务(业务侧管理入口) |
||||
|
└─ 扇出 rpcx 给 comm.ModuleConfigServices 各服务 |
||||
|
*/ |
||||
|
package svccfg |
||||
|
|
||||
|
import ( |
||||
|
"yunyan/comm" |
||||
|
"yunyan/lego/core" |
||||
|
"yunyan/modules" |
||||
|
) |
||||
|
|
||||
|
func NewModule() core.IModule { |
||||
|
return new(SvcCfg) |
||||
|
} |
||||
|
|
||||
|
type SvcCfg struct { |
||||
|
modules.ModuleBase |
||||
|
cfg *configComp |
||||
|
} |
||||
|
|
||||
|
func (this *SvcCfg) GetType() core.M_Modules { |
||||
|
return comm.ModuleSvcCfg |
||||
|
} |
||||
|
|
||||
|
func (this *SvcCfg) NewOptions() (options core.IModuleOptions) { |
||||
|
return new(modules.Options) |
||||
|
} |
||||
|
|
||||
|
func (this *SvcCfg) OnInstallComp() { |
||||
|
this.ModuleBase.OnInstallComp() |
||||
|
this.cfg = this.RegisterComp(new(configComp)).(*configComp) |
||||
|
} |
||||
@ -0,0 +1,22 @@ |
|||||
|
package main |
||||
|
|
||||
|
import ( |
||||
|
"os" |
||||
|
|
||||
|
"yunyan/comm" |
||||
|
"yunyan/lego/sys/mysql" |
||||
|
"yunyan/sys/email" |
||||
|
) |
||||
|
|
||||
|
// api 服务的业务配置热重载装配。
|
||||
|
//
|
||||
|
// api 自己持有 email 客户端(会发邮件),所以它必须能重载自己那一份——
|
||||
|
// 只重载 home 而不管 api,api 就会一直用着旧凭据。
|
||||
|
// 具体的下发/汇总入口在 modules/api/api_moduleconfig.go;本文件只登记「本服务能重载什么」。
|
||||
|
//
|
||||
|
// 注意:api **不是**配置属主,不做 seed,也不响应 reset 的实际还原(见 comm.ModuleConfigOwner)。
|
||||
|
|
||||
|
func (this *Service) setupModuleReload() { |
||||
|
comm.RegisterModuleReloader("email", func(cfg map[string]interface{}) error { return email.Reload(cfg) }) |
||||
|
comm.SetModuleReloadContext(mysql.GetSys(), os.Getenv("FIELD_ENCRYPT_KEY"), this.GetSettings().Sys) |
||||
|
} |
||||
@ -0,0 +1,33 @@ |
|||||
|
package main |
||||
|
|
||||
|
import ( |
||||
|
"os" |
||||
|
|
||||
|
"yunyan/comm" |
||||
|
"yunyan/lego/sys/mysql" |
||||
|
"yunyan/sys/email" |
||||
|
"yunyan/sys/sms" |
||||
|
) |
||||
|
|
||||
|
// 「后台改配置 → 目标服务实时生效 → 回执确认」的业务侧装配。
|
||||
|
//
|
||||
|
// 传输走 HTTP:console 是单例服务(lego/base/single)、不接 ETCD/rpcx 集群,只能从公网网关入口调进来。
|
||||
|
// 具体接口见 modules/user/api_reloadmoduleconfig.go(路由 user_reloadmoduleconfig,网关强制验签)。
|
||||
|
// 本文件只负责两件事:登记「哪些模块能热更」,以及「重载时用哪个库/哪把密钥/哪份 Sys」。
|
||||
|
|
||||
|
// registerModuleReloaders 注册可热重载的业务模块(键 = ModuleDef.SysKey)。
|
||||
|
//
|
||||
|
// 只有在此登记的模块才会真的热更;其余(各 auth / 各 pay)重载时如实回 restart_required——
|
||||
|
// 它们的 sys 尚未提供「校验失败不 panic 且保留旧实例」的 Reload,贸然复用 OnInit 会在
|
||||
|
// 配置填错时把线上进程打挂。要扩展就照 sys/email 的 Reload 补齐后再加一行。
|
||||
|
func registerModuleReloaders() { |
||||
|
comm.RegisterModuleReloader("email", func(cfg map[string]interface{}) error { return email.Reload(cfg) }) |
||||
|
comm.RegisterModuleReloader("sms", func(cfg map[string]interface{}) error { return sms.Reload(cfg) }) |
||||
|
} |
||||
|
|
||||
|
// setupModuleReload 登记重载上下文。必须在各 sys OnInit 与 mysql.OnInit 之后调用——
|
||||
|
// 重载做的是「替换已初始化好的客户端」,且要用已连上的业务库重读配置。
|
||||
|
func (this *Service) setupModuleReload() { |
||||
|
registerModuleReloaders() |
||||
|
comm.SetModuleReloadContext(mysql.GetSys(), os.Getenv("FIELD_ENCRYPT_KEY"), this.GetSettings().Sys) |
||||
|
} |
||||
@ -0,0 +1,69 @@ |
|||||
|
package email |
||||
|
|
||||
|
import "testing" |
||||
|
|
||||
|
// 热重载的核心安全属性:后台填错配置时,Reload 必须「返回 error + 保留旧客户端 + 不 panic」,
|
||||
|
// 绝不能让线上服务失去发信能力。(OnInit 走的是 panic 路径,仅用于启动。)
|
||||
|
func TestReloadBadConfigKeepsOldClient(t *testing.T) { |
||||
|
if err := OnInit(map[string]interface{}{ |
||||
|
"EmailType": 1, "FromEmail": "a@example.com", "FromName": "Old", "Password": "pw", |
||||
|
}); err != nil { |
||||
|
t.Fatalf("OnInit 失败: %v", err) |
||||
|
} |
||||
|
old := get() |
||||
|
if old == nil { |
||||
|
t.Fatal("OnInit 后客户端不应为 nil") |
||||
|
} |
||||
|
|
||||
|
cases := []struct { |
||||
|
name string |
||||
|
cfg map[string]interface{} |
||||
|
}{ |
||||
|
{"FromEmail 是裸域名(Resend 会 422 拒绝)", map[string]interface{}{ |
||||
|
"EmailType": 3, "FromEmail": "mail.example.com", "Password": "re_x"}}, |
||||
|
{"Password 为空", map[string]interface{}{ |
||||
|
"EmailType": 1, "FromEmail": "a@example.com"}}, |
||||
|
{"FromEmail 为空", map[string]interface{}{ |
||||
|
"EmailType": 1, "Password": "pw"}}, |
||||
|
} |
||||
|
for _, c := range cases { |
||||
|
t.Run(c.name, func(t *testing.T) { |
||||
|
if err := Reload(c.cfg); err == nil { |
||||
|
t.Fatal("非法配置必须返回 error") |
||||
|
} |
||||
|
if get() != old { |
||||
|
t.Fatal("重载失败后必须保留旧客户端") |
||||
|
} |
||||
|
}) |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
// 合法配置应当整体热替换客户端,且能在 SMTP ↔ Resend 两种不同具体类型间切换。
|
||||
|
func TestReloadSwapsClient(t *testing.T) { |
||||
|
if err := OnInit(map[string]interface{}{ |
||||
|
"EmailType": 1, "FromEmail": "a@example.com", "Password": "pw", |
||||
|
}); err != nil { |
||||
|
t.Fatalf("OnInit 失败: %v", err) |
||||
|
} |
||||
|
old := get() |
||||
|
if _, ok := old.(*Email); !ok { |
||||
|
t.Fatalf("EmailType=1 应为 *Email,实得 %T", old) |
||||
|
} |
||||
|
|
||||
|
if err := Reload(map[string]interface{}{ |
||||
|
"EmailType": 3, "FromEmail": "noreply@mail.example.com", "FromName": "New", "Password": "re_x", |
||||
|
}); err != nil { |
||||
|
t.Fatalf("合法配置重载应成功: %v", err) |
||||
|
} |
||||
|
cur := get() |
||||
|
if cur == old { |
||||
|
t.Fatal("重载成功后应替换为新客户端") |
||||
|
} |
||||
|
r, ok := cur.(*Resend) |
||||
|
if !ok { |
||||
|
t.Fatalf("EmailType=3 应为 *Resend,实得 %T", cur) |
||||
|
} |
||||
|
if r.ApiKey != "re_x" || r.From != "noreply@mail.example.com" || r.Name != "New" { |
||||
|
t.Fatalf("新客户端未取到新配置: %+v", r) |
||||
|
} |
||||
|
} |
||||
@ -1,25 +1,70 @@ |
|||||
package sms |
package sms |
||||
|
|
||||
|
import ( |
||||
|
"errors" |
||||
|
"sync" |
||||
|
) |
||||
|
|
||||
type ( |
type ( |
||||
ISys interface { |
ISys interface { |
||||
SendCaptcha(mobile string, captcha string) error |
SendCaptcha(mobile string, captcha string) error |
||||
} |
} |
||||
) |
) |
||||
|
|
||||
|
// defsys 由 OnInit(启动) 与 Reload(运行时热更) 写、SendCaptcha 读,故必须加锁:
|
||||
|
// 无锁并发读写接口值是 data race。
|
||||
var ( |
var ( |
||||
|
mu sync.RWMutex |
||||
defsys ISys |
defsys ISys |
||||
) |
) |
||||
|
|
||||
|
var errNotInit = errors.New("sms 系统未初始化") |
||||
|
|
||||
|
func get() ISys { |
||||
|
mu.RLock() |
||||
|
defer mu.RUnlock() |
||||
|
return defsys |
||||
|
} |
||||
|
|
||||
|
func set(sys ISys) { |
||||
|
mu.Lock() |
||||
|
defsys = sys |
||||
|
mu.Unlock() |
||||
|
} |
||||
|
|
||||
func OnInit(config map[string]interface{}, option ...Option) (err error) { |
func OnInit(config map[string]interface{}, option ...Option) (err error) { |
||||
defsys, err = newSys(newOptions(config, option...)) |
sys, err := newSys(newOptions(config, option...)) |
||||
|
if err != nil { |
||||
|
return |
||||
|
} |
||||
|
set(sys) |
||||
return |
return |
||||
} |
} |
||||
|
|
||||
|
// Reload 运行时热替换短信客户端(console 后台改配置后由业务侧调用)。
|
||||
|
// 配置不合法时返回 error 而非 panic,且**保留原客户端不动**。
|
||||
|
func Reload(config map[string]interface{}, option ...Option) error { |
||||
|
options, err := newOptionsChecked(config, option...) |
||||
|
if err != nil { |
||||
|
return err |
||||
|
} |
||||
|
sys, err := newSys(options) |
||||
|
if err != nil { |
||||
|
return err |
||||
|
} |
||||
|
set(sys) |
||||
|
return nil |
||||
|
} |
||||
|
|
||||
func NewSys(option ...Option) (sys ISys, err error) { |
func NewSys(option ...Option) (sys ISys, err error) { |
||||
sys, err = newSys(newOptionsByOption(option...)) |
sys, err = newSys(newOptionsByOption(option...)) |
||||
return |
return |
||||
} |
} |
||||
|
|
||||
func SendCaptcha(_email string, _captcha string) error { |
func SendCaptcha(_email string, _captcha string) error { |
||||
return defsys.SendCaptcha(_email, _captcha) |
sys := get() |
||||
|
if sys == nil { |
||||
|
return errNotInit |
||||
|
} |
||||
|
return sys.SendCaptcha(_email, _captcha) |
||||
} |
} |
||||
|
|||||
Loading…
Reference in new issue