You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
107 lines
5.1 KiB
107 lines
5.1 KiB
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
|
|
}
|
|
|