# gfx **Repository Path**: codedbyte/gfx ## Basic Information - **Project Name**: gfx - **Description**: No description available - **Primary Language**: Go - **License**: Apache-2.0 - **Default Branch**: main - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-06-06 - **Last Updated**: 2026-07-23 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # gfx - Go 基础框架层 gfx 是一个为 Go 项目设计的基础框架层,提供了从服务启动到可观测性、异步处理、三方适配等全链路通用能力。它帮助你快速搭建稳定、可扩展的生产级应用。 ## Templates - [`templates/service`](templates/service/README.md) 是面向生产服务的起步模板,提供基础设施接线、受保护路由和可选 `wsx.ServiceClient`,但不预置用户、订单等领域代码。 - [`templates/websocket-gateway`](templates/websocket-gateway/README.md) 是独立部署的 WebSocket Gateway,负责终端连接、首帧认证、ServiceClient 路由和集群广播,不与 Service 模板联合编排。 - [`examples/full`](examples/full/README.md) 是完整领域示例,用于演示推荐分层与 gfx 接线方式;它包含 Postgres、Redis、Notify、WebSocket、RPCX + etcd、Casbin 权限和订单数据过滤,不是最小生产模板。 ## 生产使用注意事项与能力边界 - `wsx` 提供 Fiber WebSocket 基础连接能力,以及 `Gateway` + `ServiceClient` 的客户端接入网关模型。默认客户端请求 JSON 协议只使用 `version`、`code`、`payload`;`request_id`、`server_id`、`headers`、`metadata` 会被显式拒绝,其他未知字段按标准 JSON 解码行为忽略。路由、选择器、discovery、middleware、hooks 和首帧认证器可替换扩展;浏览器端 SDK 与复杂服务治理仍由业务层负责。 - `config.Watch` 当前仅支持文件型配置源,并通过固定间隔轮询文件状态触发 reload。调用方应使用可取消的 context 管理 watcher 生命周期;`Store.Close()` 会等待后台 watcher 退出,但不会主动取消传入的 context。 - `rpx` 相关测试依赖可用 etcd。没有 etcd 时测试会 skip;本机有 etcd 时会执行集成式测试,并清理测试使用的 base path。 - `fiberx.JSON` 与 `fiberx.Error` 是项目统一响应约定。业务应保持响应 envelope 一致,避免在同一服务中混用不同响应格式。 ## 配置体系(第一版) `config` 包提供 YAML 静态加载、环境变量覆盖、动态重载和配置校验能力。 ### 一次性加载 `Load` 适合启动阶段读取配置。加载流程是: - 从 `Source` 读取原始内容,例如 `config.YAMLFile("config.yaml")`。 - 将 YAML 解码到泛型配置结构体。 - 如果配置实现了 `Validate() error`,会自动执行自校验。 - 如果传入了 `WithValidators`,会继续执行额外校验函数。 - 如果传入了 `WithEnvPrefix`,会用环境变量覆盖同名配置字段。 环境变量名按结构体字段路径生成,格式为 `_`,字段名会转为 大写蛇形命名。例如 `Server.Port` 对应 `APP_SERVER_PORT`, `Server.ShutdownTimeout` 对应 `APP_SERVER_SHUTDOWN_TIMEOUT`。当前支持覆盖 `string`、`bool`、整数、无符号整数、浮点数、`time.Duration` 和这些类型的指针; `time.Duration` 使用 Go duration 字符串,例如 `10s`、`5m`。 ```go type AppConfig struct { Server struct { Port int `yaml:"port"` } `yaml:"server"` } func (c AppConfig) Validate() error { if c.Server.Port <= 0 { return errors.New("server.port must be positive") } return nil } cfg, err := config.Load[AppConfig]( context.Background(), config.YAMLFile("config.yaml"), config.WithEnvPrefix[AppConfig]("APP"), ) if err != nil { panic(err) } ``` ### 动态监听 `Watch` 适合长期运行的服务。它会先同步执行一次 `Load`,拿到初始配置后返回 `Store[T]`,再在后台按固定间隔轮询文件修改时间。 ```go store, err := config.Watch[AppConfig]( context.Background(), config.YAMLFile("config.yaml"), config.WithPollInterval[AppConfig](time.Second), config.WithErrorHandler[AppConfig](func(err error) { log.Printf("config reload failed: %v", err) }), ) if err != nil { panic(err) } defer store.Close() ``` `Watch` 的行为需要注意: - 初始加载失败会直接返回 `error`,不会启动后台监听。 - 当前版本只支持文件型配置源,例如 `YAMLFile`。 - 默认轮询间隔是 `1s`,可以用 `WithPollInterval` 调整。 - 文件修改后会重新执行完整加载流程,包括 YAML 解码、环境变量覆盖和校验。 - 只有 reload 成功时才会更新 `Store`。 - reload 失败时会调用 `WithErrorHandler`,旧配置继续生效。 - 如果没有设置 `WithErrorHandler`,后台 reload 错误会被忽略,旧配置仍然保留。 - 调用 `store.Close()` 会关闭订阅并等待后台 watcher 退出。 ### Store 使用方式 `Store[T]` 保存当前生效配置。业务代码可以随时通过 `Get` 读取最新快照,也可以通过 `Subscribe` 监听成功更新后的配置。 ```go store, err := config.Watch[AppConfig]( context.Background(), config.YAMLFile("config.yaml"), config.WithPollInterval[AppConfig](time.Second), config.WithErrorHandler[AppConfig](func(err error) { log.Printf("config reload failed: %v", err) }), ) if err != nil { panic(err) } defer store.Close() current := store.Get() log.Printf("server port: %d", current.Server.Port) unsubscribe := store.Subscribe(func(next AppConfig) { log.Printf("config reloaded, server port: %d", next.Server.Port) }) defer unsubscribe() ``` `Store` 的语义: - `Get` 返回当前最新配置快照。 - `Subscribe` 只在配置成功 reload 并更新后触发。 - `Subscribe` 返回的函数用于取消订阅。 - 订阅回调中的 panic 会被隔离,不会影响其他订阅者或 watcher。 - `Close` 后不再接受新的订阅,已有订阅会被清空。 ### 错误处理 配置错误会包装为 `*config.Error`,其中 `Kind` 标识失败阶段: - `KindRead`:读取配置源失败。 - `KindDecode`:YAML 解码失败。 - `KindEnv`:环境变量覆盖失败。 - `KindValidation`:配置校验失败。 - `KindWatch`:监听机制失败,例如非文件源不支持 watch。 可以用 `errors.As` 判断错误类型: ```go var cfgErr *config.Error if errors.As(err, &cfgErr) { log.Printf("config failed, kind=%s source=%s op=%s", cfgErr.Kind, cfgErr.Source, cfgErr.Op) } ``` ## Middleware 限流与幂等 `fiberx` 提供限流和幂等中间件。默认内置进程内内存实现,适合单实例服务、开发环境和测试场景。多副本生产部署可以通过 `redisx` 注入 Redis 实现。 ### 限流配置 ```yaml middleware: rateLimit: enabled: true rate: 10 burst: 20 keyBy: ip ttl: 5m cleanupInterval: 1m ``` `rate` 表示每秒补充 token 数,`burst` 表示桶容量。`keyBy` 支持 `ip`、`path` 和 `header:`。 ```go handler, err := fiberx.RateLimitFromConfig(cfg.Middleware.RateLimit) if err != nil { panic(err) } app.Use(handler) ``` 超限请求返回统一 `429` 响应,并回写 `X-RateLimit-Limit`、`X-RateLimit-Remaining`、`X-RateLimit-Reset` 和 `Retry-After`。 ### 幂等配置 ```yaml middleware: idempotency: enabled: true headerName: Idempotency-Key methods: [POST, PUT, PATCH, DELETE] requireKey: false ttl: 24h maxBodyBytes: 1048576 ``` 默认缺少 `Idempotency-Key` 时放行。设置 `requireKey: true` 后,缺少 key 的写请求返回 `400`。 ```go handler, err := fiberx.IdempotencyFromConfig(cfg.Middleware.Idempotency) if err != nil { panic(err) } app.Use(handler) ``` 同一个 key 的请求正在处理时返回 `409`。首次请求完成后,重复请求会回放已缓存响应,并回写 `Idempotency-Replayed: true`。`5xx` 响应和超过 `maxBodyBytes` 的响应不会缓存。 ## JWT 鉴权 `pkg/jwt` 提供框架无关的 JWT 签发、解析、校验与 access/refresh 双 token 管理能力。`jwtx` 负责 `Config` 聚合、`Manager` 工厂和 `server.Initializer` 注册;`fiberx` 提供 Bearer 鉴权中间件。 默认 refresh 模式为无状态双 JWT;生产环境可切换 opaque refresh,并通过 `redisx.NewRefreshStore` 注入 `jwt.RefreshStore`。 默认 TTL:access token `15m`,refresh token `7d`(`168h`)。 ### JWT 配置 ```yaml jwt: algorithm: HS256 secret: ${JWT_SECRET} issuer: my-app accessTTL: 15m refreshTTL: 168h refreshMode: jwt middleware: skipPaths: - /healthz - /readyz - /api/auth/login - /api/auth/refresh auth: enabled: true ``` `refreshMode` 支持 `jwt`(默认,无状态双 JWT)和 `opaque`(随机 refresh token 存 Redis)。 `middleware.skipPaths` 按请求 `path` 精确匹配跳过鉴权和授权,语义与 `WithAccessLogSkipPaths` 一致。 ```go type AppConfig struct { Server server.Config `yaml:"server"` JWT jwt.Config `yaml:"jwt"` Middleware struct { Auth fiberx.AuthConfig `yaml:"auth"` SkipPaths []string `yaml:"skipPaths"` } `yaml:"middleware"` } func (c AppConfig) ServerConfig() server.Config { return c.Server } ``` ### Manager 注册 `jwtx.Initializer` 向 `server.Container` 注册 `jwt.Manager`,默认依赖名为 `jwt`。 ```go app, err := server.Load[AppConfig]( context.Background(), config.YAMLFile("config.yaml"), server.WithServerOptions[AppConfig]( server.WithInitializer(jwtx.Initializer[AppConfig]("", func(c AppConfig) jwtx.Config { return jwtx.Config{JWT: c.JWT} })), ), ) if err != nil { panic(err) } ``` ### 鉴权中间件 `fiberx.AuthFromConfig` 从 YAML 配置创建中间件;`enabled: false` 时返回 no-op handler。 `fiberx.AuthRequired` 适合不依赖 YAML 开关、始终要求鉴权的场景。两者都需要 `jwt.Parser`; `jwt.Manager` 不实现 `Parser`,业务可通过 `jwt.NewParser(cfg.JWT)` 创建 parser,或在 initializer 中单独注册。 ```go parser, err := jwt.NewParser(cfg.JWT) if err != nil { panic(err) } handler, err := fiberx.AuthFromConfig(cfg.Middleware.Auth, parser) if err != nil { panic(err) } app.Use(handler) ``` ```go app.Use(fiberx.AuthRequired(parser)) app.Get("/me", func(c fiber.Ctx) error { claims, ok := fiberx.Claims(c) if !ok { return fiberx.Error(c, result.NewError(result.CodeUnauthorized, result.MessageUnauthorized)) } return fiberx.JSON(c, map[string]string{"sub": claims.Subject}) }) ``` token 缺失或无效时返回统一 `401` 响应。中间件只接受 `Authorization: Bearer ` 格式的 access token,校验通过后通过 `fiberx.Claims` 或 `fiberx.TypedClaims[T]` 读取 claims。 自定义 claims 字段必须带 `json` tag,`ParseClaims[T]` 与 `TypedClaims` 通过 JSON round-trip 实现类型转换: ```go type AppClaims struct { jwt.RegisteredClaims UserID string `json:"uid"` Roles []string `json:"roles"` } claims, ok := fiberx.TypedClaims[AppClaims](c) ``` ### Opaque refresh 模式 生产环境可将 `refreshMode` 设为 `opaque`,refresh token 为随机字符串并存入 Redis。 initializer 注册顺序:`redis` → `jwt-refresh-store` → `jwt`(向 `jwtx.Config.RefreshStore` 注入 store)。 ```yaml jwt: algorithm: HS256 secret: ${JWT_SECRET} refreshMode: opaque ``` ```go server.WithInitializer(redisx.Initializer[AppConfig]("", func(cfg AppConfig) redisx.Config { return cfg.Redis })), server.WithInitializerFunc[AppConfig]("jwt-refresh-store", func(ctx context.Context, cfg AppConfig, deps *server.Container) error { if cfg.JWT.RefreshMode != "opaque" { return nil } client, ok := server.Get[*redisx.Client](deps, redisx.DEPS_NAME) if !ok { return errors.New("redis dependency is required for opaque refresh mode") } store, err := redisx.NewRefreshStore(client, redisx.RefreshStoreConfig{}) if err != nil { return err } return deps.Set("jwtRefreshStore", store) }), server.WithInitializerFunc[AppConfig]("jwt", func(ctx context.Context, cfg AppConfig, deps *server.Container) error { jwtxCfg := jwtx.Config{JWT: cfg.JWT} if cfg.JWT.RefreshMode == "opaque" { store, ok := server.Get[*redisx.RefreshStore](deps, "jwtRefreshStore") if !ok { return errors.New("jwt refresh store is required for opaque mode") } jwtxCfg.RefreshStore = store } mgr, err := jwtx.NewManager(jwtxCfg) if err != nil { return err } return deps.Set(jwtx.DEPS_NAME, mgr) }), ``` ### 业务端点 `examples/full` 不演示 login / refresh 端点;认证流程由业务项目自行实现,通过 `jwt.Manager` 的 `IssueTokenPair`、`RefreshAccess` 和 `RevokeRefresh` 完成 token 生命周期管理。 ```go mgr, ok := server.Get[jwt.Manager](deps, jwtx.DEPS_NAME) if !ok { return fiberx.Error(c, result.NewError(fiber.StatusInternalServerError, "jwt unavailable")) } pair, err := mgr.IssueTokenPair(c.Context(), jwt.RegisteredClaims{Subject: userID}) if err != nil { return err } return fiberx.JSON(c, pair) ``` ```go pair, err := mgr.RefreshAccess(c.Context(), req.RefreshToken) if err != nil { return fiberx.Error(c, result.NewError(result.CodeUnauthorized, result.MessageUnauthorized)) } return fiberx.JSON(c, pair) ``` opaque 模式下 `RevokeRefresh` 会从 Redis 删除 refresh token;jwt 模式下为 no-op,由客户端丢弃 token 即可。 ## 多租户接口权限 `pkg/authz` 定义框架无关的授权接口和上下文类型,`casbinx` 将 Casbin enforcer 适配为 `authz.Authorizer` 并注册到 `server.Container`,`fiberx.AuthorizeFromConfig` 在 JWT 鉴权之后执行接口权限校验。`fiberx` 只依赖 `authz.Authorizer`,不直接依赖 Casbin。 权限检查使用四元组: ```go type Request struct { Subject string TenantID string Object string Action string } ``` `Subject` 来自 `fiberx.Claims(c).Subject`,`TenantID` 默认来自 `X-Tenant-ID` 请求头。 校验通过后,中间件会写入 `authz.Context{Subject, TenantID}`,handler 或 repository 可通过 `fiberx.AuthzContext(c)` 或 `authz.FromContext(c.Context())` 读取。`authz.Context` 不包含角色; 角色解析和数据权限解析属于业务层职责。 ### Casbin 配置与注册 `casbinx` 依赖已注册的 `postgresx.Client`,默认从 `casbin_rule` 表加载 policy。内置 model 是 RBAC with domains:`sub, dom, obj, act`。`dom` 对应租户 ID。 ```yaml casbin: enabled: true tableName: casbin_rule modelFile: model.conf autoLoadInterval: 30s middleware: skipPaths: - /healthz - /readyz permission: enabled: true ``` `autoLoadInterval` 为 `0` 或不填时只在启动时加载 policy;大于 `0` 时后台定时执行 `LoadPolicy`,并在服务关闭时通过 `server.Container` 的 closer 停止 goroutine。 `modelFile` 指向 Casbin model 文件,默认读取当前工作目录下的 `model.conf`;更复杂场景 也可以用 `casbinx.WithModel` 注入运行时构造的 Casbin model。 ```go type AppConfig struct { Server server.Config `yaml:"server"` Postgres postgresx.Config `yaml:"postgres"` Casbin casbinx.Config `yaml:"casbin"` Middleware struct { Auth fiberx.AuthConfig `yaml:"auth"` Permission fiberx.AuthorizeConfig `yaml:"permission"` SkipPaths []string `yaml:"skipPaths"` } `yaml:"middleware"` } cfg, err := config.Load[AppConfig](context.Background(), config.YAMLFile("config.yaml")) if err != nil { panic(err) } app, err := server.New( cfg, server.WithInitializer(postgresx.Initializer[AppConfig]("", func(c AppConfig) postgresx.Config { return c.Postgres })), server.WithInitializer(casbinx.Initializer[AppConfig]("", func(c AppConfig) casbinx.Config { return c.Casbin })), ) ``` ### 授权中间件 推荐中间件顺序是:`fiberx.AuthFromConfig` 或 `fiberx.AuthRequired` 先解析 JWT, 再挂载 `fiberx.AuthorizeFromConfig`。授权中间件缺少 claims 时返回 `401`,缺少租户头或 Authorizer 拒绝时返回统一 `403`。 默认资源映射是请求 path + HTTP method,例如 `GET /api/orders/42` 映射为 `Object="/api/orders/42", Action="GET"`。生产系统通常应通过 `WithResourceMapper` 显式映射为稳定业务资源,避免把动态 ID 写入 policy。 下面的 mapper 只是示例;实际项目应按路由表或业务资源命名维护稳定映射。 ```go server.WithRoutes[AppConfig](func(app *fiber.App, deps *server.Container) error { parser, err := jwt.NewParser(cfg.JWT) if err != nil { return err } authn, err := fiberx.AuthFromConfig(cfg.Middleware.Auth, parser) if err != nil { return err } authorizer, ok := server.Get[authz.Authorizer](deps, casbinx.DEPS_NAME) if !ok { return errors.New("casbin authorizer unavailable") } authzHandler, err := fiberx.AuthorizeFromConfig( cfg.Middleware.Permission, authorizer, fiberx.WithResourceMapper(func(c fiber.Ctx) (string, string) { if strings.HasPrefix(c.Path(), "/api/orders/") { return "/api/orders/:id", c.Method() } return c.Path(), c.Method() }), ) if err != nil { return err } app.Use(authn) app.Use(authzHandler) return nil }) ``` Casbin policy 示例: ```text p, admin, tenant-a, /api/orders/:id, GET g, user-1, admin, tenant-a ``` ### 数据权限类型 `pkg/authz` 还提供 `DataScope` 与 `DataFilter` 作为稳定 DTO,用于业务 repository 表达数据过滤结果。它不包含 SQL,也不定义 `DataScopeResolver`;不同业务的组织结构、角色来源和 部门继承规则不同,应在业务项目中实现。 ## WebSocket 服务端基础能力 `wsx` 包基于 Fiber WebSocket 提供底层升级、连接生命周期、心跳、单播和广播能力。`Hub`/`Server` 仍可用于简单的进程内 WebSocket 服务: ```go hub, err := wsx.NewHub(wsx.Config{ ReadLimit: 1 << 20, SendQueueSize: 64, }) if err != nil { panic(err) } srv, err := wsx.NewServer(hub, wsx.HandlerFunc(func(ctx context.Context, conn *wsx.Conn, msg wsx.Message) error { return conn.Send(msg) })) if err != nil { panic(err) } app.Get("/ws", srv.FiberHandler()) ``` 可以通过 `hub.Send(connID, msg)` 向指定连接发送消息,也可以通过 `hub.Broadcast(msg)` 向当前在线连接广播消息。默认情况下,新连接使用已有 `connID` 注册时会替换并关闭旧连接,`Hub` 中只保留最新连接。服务关闭时调用 `hub.Close(ctx)` 可关闭所有连接。 ### Gateway 接入网关 `Gateway` 现在是客户端 WebSocket 接入网关:客户端连接到 `ClientHandler()`,业务服务通过 `ServiceClient` 注册自己的 `server_id` 和可承担的 responsibilities。Gateway 使用可替换 `Router` 把客户端 `code` 解析成 responsibilities,再通过 discovery 和 selector 选择拥有这些 responsibilities 的实例。 ```go router := wsx.NewStaticRouter(map[string][]string{ "chat.send": {"chat.message"}, }) gateway, err := wsx.NewGateway(wsx.GatewayConfig{ ServerID: "gateway-1", ClientPath: "/ws", UpstreamPath: "/_wsx/upstream", }, wsx.WithGatewayRouter(router)) if err != nil { panic(err) } app.Get("/ws", gateway.ClientHandler()) app.Get("/_wsx/upstream", gateway.UpstreamHandler()) if err := gateway.Start(ctx); err != nil { panic(err) } defer gateway.Close(ctx) ``` 默认客户端 JSON 请求只暴露业务字段,不允许客户端设置内部路由字段: ```json {"version":"wsx.v1","code":"chat.send","payload":{"text":"hello"}} ``` 默认客户端响应只包含 `code`、`payload`、`error`,不会回传内部 `version` 或 `request_id`。`server_id`、`headers`、`metadata`、`request_id` 属于 gateway 与 upstream 内部协议;如果业务需要客户端自带这些字段,应通过自定义 `Codec` 或后续协议版本扩展。 #### 首帧认证与可信身份 启用 `GatewayConfig.ClientAuth` 后,客户端连接先处于 pending 状态,只能在超时前发送配置的认证 code(默认 `auth`)。Gateway 忽略查询参数中的 `client_id`,并默认使用密码学随机数为连接生成不可预测的内部 ID。认证首帧示例: ```json {"version":"wsx.v1","code":"auth","payload":{"access_token":""}} ``` 认证器通过 `WithGatewayClientAuthenticator` 注入。认证成功响应使用同一个 code,并返回服务端生成的 `client_id`: ```json {"code":"auth","payload":{"client_id":""}} ``` `WithGatewayIDGenerator` 是兼容 option,会同时覆盖内部 request ID 与 client ID 生成器;`WithGatewayClientIDGenerator` 只覆盖 client ID 生成器。生产环境自定义任一生成器并使其参与 client ID 生成时,必须保证结果唯一且不可预测。 认证器返回的 `ClientIdentity.Subject` 必须非空。Gateway 会复制可信 metadata,并强制用 Subject 覆盖保留的 `metadata["subject"]`;后续每条 `GatewayRequest` 都携带这份服务端身份。客户端协议禁止提交内部 `metadata`,不能伪造或覆盖可信 Subject。只有 authenticated 连接可以路由业务消息、接收 Push 或 Broadcast。 首帧认证相关错误使用下列稳定错误码。Gateway 会先等待错误帧写出,再关闭连接: - `auth_required`:pending 连接发送了非认证 code,包括认证前发送业务消息。 - `auth_failed`:认证 payload 无效、认证器拒绝或异常,或者认证器返回空 Subject 等无效身份。 - `auth_timeout`:连接未在认证期限内成功完成认证;该错误标记为可重试。 - `already_authenticated`:authenticated 连接再次发送认证 code。 #### Push 与集群 Broadcast `ServiceClient.Push` 为命令生成 request ID,并等待 Gateway 返回 ACK 或稳定错误。ACK 只表示消息已进入目标 authenticated 连接的有界发送队列,不表示终端已经收到或处理;目标不存在、连接已关闭、队列已满或等待超时都会返回错误。 `ServiceClient.Broadcast` 将广播交给当前 Gateway 发布到配置的 `notify.Notifier`,发布成功后返回 ACK。每个 Gateway 使用由稳定 `server_id` 派生的独立消费组,因此所有在线 Gateway 都能消费广播;订阅新消费组时从最新消息开始。Gateway 对事件执行 TTL 检查和短期去重,再向本机 authenticated 连接非阻塞投递。默认事件 TTL 为 30 秒,去重保留时间至少为事件 TTL 的两倍;广播 ACK 不保证每个客户端已经完成网络写入。 配置集群广播时使用 `WithGatewayNotifier`,并确保每个 Gateway 的 `server_id` 在部署中稳定且唯一。Redis Streams 和 NATS JetStream 都保留同一消费组内的竞争消费语义,而 Gateway 间使用不同 group 获得全节点广播语义。 #### Upstream 网络边界 `/_wsx/upstream` 是 `ServiceClient` 的内部 WebSocket 入口,不提供应用层鉴权。生产环境必须只允许可信内网访问,并通过防火墙、Kubernetes NetworkPolicy 或等价网络策略阻止公网和非可信工作负载连接。 业务服务侧使用 `ServiceClient`。业务启动配置声明该进程承担的 responsibilities,action/code 列表由业务代码自己的注册表维护,Gateway 的 route table 负责把 `code -> responsibilities`: ```yaml ws: enabled: true service: serverId: chat-1 responsibilities: - chat.message ``` ```go client, err := wsx.NewServiceClient(wsx.ServiceClientConfig{ ServerID: cfg.WS.Service.ServerID, Responsibilities: cfg.WS.Service.Responsibilities, }, wsx.ServiceHandlerFunc(func(ctx context.Context, req wsx.UpstreamRequest) (wsx.UpstreamResponse, error) { return wsx.UpstreamResponse{Payload: json.RawMessage(`{"ok":true}`)}, nil }), wsx.WithServiceClientGatewayURL("ws://127.0.0.1:8080/_wsx/upstream"), ) if err != nil { panic(err) } if err := client.Start(ctx); err != nil { panic(err) } defer client.Close(ctx) ``` Sticky 绑定由业务响应显式建立:只有 upstream response 的 `route.sticky_server_id` 非空时,Gateway 才会记录 `client_id + responsibility -> server_id`,后续同一客户端同一 responsibility 优先使用该实例。成功响应不会自动创建 sticky 绑定。 扩展点: - `Codec`:替换默认 JSON 客户端协议和 upstream 协议。 - `Router`:替换静态 `code -> responsibilities` 路由表。 - `Selector`:替换 sticky、least-load、random 组合选择策略。 - `Discovery`:使用内存 discovery 或 etcd discovery 注册、list、watch gateway 实例;service client 通过 discovery 找到 gateway 后,在 upstream 注册帧中声明自己的 responsibilities。 - `GatewayMiddleware` 和 `GatewayHooks`:在 decode、route、select、forward、response 和连接生命周期注入鉴权、审计、指标与错误处理。 完整示例使用 `etcdx` discovery 暴露 wsx gateway upstream endpoint,service client 连接 gateway 后用 upstream 注册帧声明 responsibilities。示例配置声明 startup responsibilities;业务 action 仍在 `WSService` 代码中注册,router 在启动时从 `WSService.Codes()` 生成 `code -> cfg.WS.Service.Responsibilities` 路由表。 ## Redis 分布式基础能力 `redisx` 包基于 `github.com/redis/go-redis/v9` 提供单机 Redis client、Redis 限流器、Redis 幂等存储和 Redis 分布式锁。第一版只支持单机 Redis;Sentinel 和 Cluster 不在本版范围内。 ### Redis 配置 ```yaml redis: addr: 127.0.0.1:6379 username: "" password: "" db: 0 dialTimeout: 5s readTimeout: 3s writeTimeout: 3s keyPrefix: gfx ``` ```go type AppConfig struct { Server server.Config `yaml:"server"` Redis redisx.Config `yaml:"redis"` RateLimit fiberx.RateLimitConfig `yaml:"rateLimit"` Idempotent fiberx.IdempotencyConfig `yaml:"idempotency"` } func (c AppConfig) ServerConfig() server.Config { return c.Server } ``` ### Client 生命周期 `redisx.Client` 实现 `HealthCheck(context.Context) error` 和 `Close(context.Context) error`,注册到 `server.Container` 后会自动参与 `/readyz` 和优雅关闭。 ```go app, err := server.Load[AppConfig]( context.Background(), config.YAMLFile("config.yaml"), server.WithServerOptions[AppConfig]( server.WithInitializer(redisx.Initializer[AppConfig]("", func(cfg AppConfig) redisx.Config { return cfg.Redis })), ), ) if err != nil { panic(err) } ``` ### Redis 限流 `redisx.NewRateLimiter` 实现 `pkg/ratelimit.Limiter`,通过 `fiberx.WithRateLimiter` 注入中间件。 ```go server.WithInitializerFunc[AppConfig]("rate-limiter", func(ctx context.Context, cfg AppConfig, deps *server.Container) error { client, ok := server.Get[*redisx.Client](deps, "redis") if !ok { return errors.New("redis dependency is required") } limiter, err := redisx.NewRateLimiter(client, redisx.RateLimitConfig{ Rate: cfg.RateLimit.Rate, Burst: cfg.RateLimit.Burst, TTL: cfg.RateLimit.TTL, }) if err != nil { return err } return deps.Set("rateLimiter", limiter) }) server.WithRoutes[AppConfig](func(app *fiber.App, deps *server.Container) error { limiter, _ := server.Get[*redisx.RateLimiter](deps, "rateLimiter") handler, err := fiberx.RateLimitFromConfig(rateLimitConfig, fiberx.WithRateLimiter(limiter)) if err != nil { return err } app.Use(handler) return nil }) ``` ### Redis 幂等 `redisx.NewIdempotencyStore` 实现 `pkg/idempotency.Store`,通过 `fiberx.WithIdempotencyStore` 注入中间件。 ```go server.WithInitializerFunc[AppConfig]("idempotency-store", func(ctx context.Context, cfg AppConfig, deps *server.Container) error { client, ok := server.Get[*redisx.Client](deps, "redis") if !ok { return errors.New("redis dependency is required") } store, err := redisx.NewIdempotencyStore(client, redisx.IdempotencyConfig{}) if err != nil { return err } return deps.Set("idempotencyStore", store) }) server.WithRoutes[AppConfig](func(app *fiber.App, deps *server.Container) error { store, _ := server.Get[*redisx.IdempotencyStore](deps, "idempotencyStore") handler, err := fiberx.IdempotencyFromConfig(idempotencyConfig, fiberx.WithIdempotencyStore(store)) if err != nil { return err } app.Use(handler) return nil }) ``` ### Redis 分布式锁 `redisx.NewLocker` 提供 `TryLock`、`Unlock` 和 `WithLock`。释放锁时会比较 token,避免误删其他持有者的锁。 ```go locker, err := redisx.NewLocker(redisClient, redisx.LockConfig{TTL: 30 * time.Second}) if err != nil { panic(err) } err = locker.WithLock(context.Background(), "jobs:daily-report", 30*time.Second, func(ctx context.Context) error { return runDailyReport(ctx) }) if errors.Is(err, redisx.ErrLockNotAcquired) { return } if err != nil { panic(err) } ``` ### 异步通知 `pkg/notify` 定义中间件无关的发布、订阅、消息和 handler 抽象。`redisx.NewStreamNotifier` 和 `natsx.NewJetStreamNotifier` 分别提供基于 Redis Streams 与 NATS JetStream 的可靠实现, 均支持消费组、自动 ACK、失败重试和 DLQ。业务层应依赖 `notify.Notifier` 接口,按部署环境 注入具体实现。 Redis Streams 示例: ```go notifier, err := redisx.NewStreamNotifier(redisClient, redisx.NotifyConfig{ BatchSize: 16, PendingMinIdle: 30 * time.Second, MaxAttempts: 3, }) if err != nil { panic(err) } _, err = notifier.Publish(ctx, notify.Topic("user.created"), notify.PublishMessage{ Key: "user-1", Payload: []byte(`{"id":"user-1"}`), Headers: map[string]string{"type": "user.created"}, }) if err != nil { panic(err) } sub, err := notifier.Subscribe(ctx, notify.SubscriptionOptions{ Topic: notify.Topic("user.created"), Group: notify.Group("audit"), Consumer: notify.Consumer("api-1"), }, notify.HandlerFunc(func(ctx context.Context, msg notify.Message) error { // handler 返回 nil 后自动 ACK;返回 error 后按配置重试。 return nil })) if err != nil { panic(err) } defer sub.Close(context.Background()) ``` NATS JetStream 示例: ```go natsClient, err := natsx.NewClient(natsx.Config{ URL: "nats://127.0.0.1:4222", }) if err != nil { panic(err) } defer natsClient.Close(context.Background()) notifier, err := natsx.NewJetStreamNotifier(natsClient, natsx.NotifyConfig{ AckWait: 30 * time.Second, MaxDeliver: 3, }) if err != nil { panic(err) } ``` 同一个 topic + group 下的多个 consumer 会竞争消费,一条消息只会被其中一个 consumer 处理。相同 topic 下的不同 group 会各自消费一份消息。Redis 使用 Streams, NATS 使用 JetStream,均不是 Core Pub/Sub;后续还可以扩展 RabbitMQ、Kafka 等适配。 ## PostgreSQL 数据库基础能力 `postgresx` 基于 `pgxpool` 提供 PostgreSQL 连接池、健康检查、关闭和事务 helper。 它支持 primary + replicas 读写分离;如果 `replicas` 为空,则不启用读写分离, `Reader()`、`Replica()` 和只读事务都会使用 primary。 ```yaml postgres: primary: dsn: postgres://user:pass@primary:5432/app?sslmode=disable maxConns: 20 minConns: 2 replicas: - dsn: postgres://user:pass@replica-1:5432/app?sslmode=disable maxConns: 20 minConns: 2 healthTimeout: 2s replicaStrategy: round_robin requireReplicasHealthy: true ``` ```go type AppConfig struct { Server server.Config `yaml:"server"` Postgres postgresx.Config `yaml:"postgres"` } func (c AppConfig) ServerConfig() server.Config { return c.Server } app, err := server.Load[AppConfig]( context.Background(), config.YAMLFile("config.yaml"), server.WithServerOptions[AppConfig]( server.WithInitializer(postgresx.Initializer[AppConfig]("", func(cfg AppConfig) postgresx.Config { return cfg.Postgres })), server.WithRoutes[AppConfig](func(app *fiber.App, deps *server.Container) error { app.Get("/users", func(c fiber.Ctx) error { db, ok := server.Get[*postgresx.Client](deps, postgresx.DEPS_NAME) if !ok { return fiberx.Error(c, result.NewError(fiber.StatusInternalServerError, "database unavailable")) } rows, err := db.Reader().Query(c.Context(), "select id, name from users") if err != nil { return err } defer rows.Close() return fiberx.JSON(c, []string{}) }) return nil }), ), ) ``` 写操作或强一致读应显式使用 `Writer()`: ```go _, err := db.Writer().Exec(ctx, "insert into users(name) values($1)", name) ``` 事务默认走 primary: ```go err := postgresx.WithTx(ctx, db, func(ctx context.Context, tx pgx.Tx) error { _, err := tx.Exec(ctx, "insert into users(name) values($1)", name) return err }) ``` 只读事务走 `Reader()`;没有 replicas 时自动 fallback 到 primary: ```go err := postgresx.WithReadTx(ctx, db, func(ctx context.Context, tx pgx.Tx) error { return tx.QueryRow(ctx, "select count(*) from users").Scan(&count) }) ``` `postgresx.Client` 实现了 `HealthCheck(context.Context) error` 和 `Close(context.Context) error`,注册到 `server.Container` 后会自动参与 `/readyz` 和优雅关闭。 ### 数据库 Migration(goosex) `goosex` 基于 [goose](https://github.com/pressly/goose) 提供 PostgreSQL schema 版本化与迁移执行。 它只使用 primary DSN,通过独立短连接执行 migration,不占用 `postgresx` 连接池。 `postgresx` 负责连接池与查询;`goosex` 负责 schema 演进,两者职责互补。 ```yaml migration: enabled: true autoMigrate: false migrationsDir: db/migrations ``` - `enabled: false`:Initializer 为 no-op,不校验 `migrationsDir`。 - `enabled: true, autoMigrate: false`:Initializer 为 no-op;迁移由 CLI / Makefile / CI 执行。 - `enabled: true, autoMigrate: true`:服务启动时对 primary 自动执行 `goose up`;失败时 Initializer 返回 error,服务不启动(fail-fast)。 DSN 不单独配置。Initializer 通过注入函数从业务配置读取 `postgresx.Config.Primary.DSN`。 ```go type AppConfig struct { Server server.Config `yaml:"server"` Postgres postgresx.Config `yaml:"postgres"` Migration goosex.Config `yaml:"migration"` } func (c AppConfig) ServerConfig() server.Config { return c.Server } app, err := server.Load[AppConfig]( context.Background(), config.YAMLFile("config.yaml"), server.WithServerOptions[AppConfig]( server.WithInitializer(postgresx.Initializer[AppConfig]("", func(cfg AppConfig) postgresx.Config { return cfg.Postgres })), server.WithInitializer(goosex.Initializer[AppConfig]( "", func(cfg AppConfig) goosex.Config { return cfg.Migration }, func(cfg AppConfig) string { return cfg.Postgres.Primary.DSN }, )), ), ) ``` Initializer 应注册在 `postgresx` 之后(需要 DSN),在业务路由之前。默认依赖名为 `migration`。 CLI / CI 可复用同一套底层逻辑: ```go if err := goosex.RunUp(ctx, dsn, "db/migrations"); err != nil { panic(err) } ``` `RunDown` 与 `RunStatus` 同样可用;生产环境推荐 `autoMigrate: false`,由部署流水线显式执行 `migrate-up`。 `goosex` 当前是 PostgreSQL-only 适配包。由于 goose 的 dialect 是进程级全局状态, `goosex` 会在包内串行化 `SetDialect` 与 migration 执行窗口,避免同进程并发迁移互相影响。 ### sqlc 类型安全查询 [sqlc](https://sqlc.dev/) 是开发时工具,从 `.sql` 生成类型安全 Go 查询代码。生成代码归属 业务项目,不属于 `gfx` 库公共 API。`postgresx.QuerierFactory` 将 sqlc 生成的 `Queries` 绑定到 `Writer()` / `Reader()` / `WithTx`,与读写分离约定一致。 推荐目录结构(以 `examples/full` 为参考): ```text db/ migrations/ # goose SQL(带 -- +goose Up/Down 注解) 00001_init.sql schema/ # sqlc 按聚合拆分的表定义(与 migrations 保持同步) users.sql orders.sql queries/ # sqlc 查询 SQL(按聚合拆分) users.sql orders.sql sqlc.yaml # 多个 sql block,每聚合一个生成包 internal/ storage/ userdb/ # sqlc 生成(package userdb) orderdb/ user.go # Repository,映射 userdb → model ``` `sqlc.yaml` 按聚合拆多个 `sql` block,避免所有 model 与查询方法挤在同一包: ```yaml version: "2" sql: - engine: postgresql queries: db/queries/users.sql schema: db/schema/users.sql gen: go: package: userdb out: internal/storage/userdb sql_package: pgx/v5 - engine: postgresql queries: db/queries/orders.sql schema: db/schema/orders.sql gen: go: package: orderdb out: internal/storage/orderdb sql_package: pgx/v5 ``` - **goose** 的 schema 单一来源为 `db/migrations`;**sqlc** 的 `db/schema/*.sql` 从 migrations 提取各表 DDL,修改表结构时需同步更新。 - `sql_package: pgx/v5` 使生成代码兼容 `postgresx.DBTX`(`*pgxpool.Pool` 与 `pgx.Tx` 均满足)。 - sqlc 生成的 row struct(如 `userdb.User`)是持久化 DTO,与业务 `internal/model` 领域实体分离。 - 查询参数优先使用 `sqlc.arg(name)` 等命名参数;必要时显式 cast,避免生成 `ColumnN` 或 `interface{}` 字段削弱类型边界。 - 修改 schema 或 query 后必须重新运行 `sqlc generate`,并确认生成代码已提交且重复生成不产生 diff。 `QuerierFactory` 用法(每个聚合独立绑定): ```go factory, err := postgresx.NewQuerierFactory(client, userdb.New) if err != nil { panic(err) } // 写操作或强一致读 user, err := factory.Writer().CreateUser(ctx, name) // 只读查询(有 replica 时走 replica) users, err := factory.Reader().ListUsers(ctx, limit) // 事务 err = factory.WithTx(ctx, func(ctx context.Context, q *userdb.Queries) error { _, err := q.CreateUser(ctx, name) return err }) ``` 约定:写操作与强一致读使用 `Writer()`;可接受延迟的只读查询使用 `Reader()`;事务内通过 `WithTx` 获取绑定到 `pgx.Tx` 的 `Queries` 实例。不同聚合使用各自的 `QuerierFactory`。 ### 请求日志 `fiberx.AccessLog` 提供 handler 请求完成后的结构化日志。它复用 `pkg/logger.Logger`,由业务通过 `server.WithMiddleware(...)` 显式挂载,不会由 `server.New` 默认启用。 业务层 logger 注入职责划分可参考 [`examples/full/README.md`](examples/full/README.md) 的“日志规范”小节。 ```go log, err := logger.New(serverCfg.Log, logger.WithServiceName(serverCfg.Name)) if err != nil { panic(err) } app, err := server.New(cfg, server.WithLogger[AppConfig](log), server.WithMiddleware[AppConfig]( fiberx.AccessLog(log), ), ) ``` 日志消息固定为 `request completed`。默认字段包括 `method`、`url`、`path`、`route`、`headers`、`query`、`params`、`body`、`body_truncated`、`status`、`response`、`response_truncated`、`latency_ms`、`ip` 和 `error`;启用 `fiberx.Trace` 时还会包含 `requestId` 和 `traceId`。第一版不会对 header、请求体或响应体做脱敏;请求体和响应体默认最多记录 `4096` 字节,只有发生截断时才记录对应的 `_truncated` 字段。 日志级别根据 handler 完成后的结果选择: - `Info`:响应状态码小于 `400` 且 handler 没有返回错误。 - `Warn`:响应状态码为 `4xx`,包括返回 `*fiber.Error` 且 code 为 `4xx` 的场景。 - `Error`:响应状态码为 `5xx`、handler 返回普通错误,或 panic 被上游 `fiberx.Recover` 转成错误响应的场景。 如果 handler 返回普通错误且当前响应状态码仍小于 `400`,请求日志会按 `500` 记录 `status`;返回 `*fiber.Error` 时会优先使用其中的 HTTP code。`log == nil` 时中间件直接透传 `c.Next()`,不记录日志。 ```go fiberx.AccessLog( log, fiberx.WithAccessLogMaxBodyBytes(8192), fiberx.WithAccessLogSkipPaths("/healthz", "/readyz"), ) ``` `WithAccessLogSkipPaths` 使用请求 path 精确匹配,命中后不记录请求日志。`WithAccessLogMaxBodyBytes(0)` 使用默认值,传入负数会关闭请求体和响应体记录。 ### 可观测性(Metrics + Tracing) `pkg/telemetry` 与 `otelx` 提供基于 OpenTelemetry 的 HTTP metrics 与 trace 导出能力。 默认 `server.observability.enabled: false`,不影响现有服务。启用后经 OTLP 导出到 Collector(可扇出到 Prometheus / Jaeger / Tempo)。 ```yaml server: name: app observability: enabled: true serviceVersion: "1.0.0" otlp: endpoint: otel-collector:4317 protocol: grpc insecure: true trace: enabled: true sampleRatio: 1.0 metrics: enabled: true interval: 15s prometheus: enabled: false path: /metrics ``` 接线示例: ```go app, err := server.Load[AppConfig]( ctx, config.YAMLFile("config.yaml"), server.WithServerOptions[AppConfig]( server.WithInitializer(otelx.Initializer[AppConfig]("", func(c AppConfig) otelx.Config { return c.Server.Observability }, func(c AppConfig) string { return c.Server.Name }, )), server.WithRoutes[AppConfig]( otelx.PrometheusRoutes[AppConfig](), ), server.WithMiddleware[AppConfig]( fiberx.Otel(), fiberx.AccessLog(log), ), ), ) ``` 推荐中间件顺序(外 → 内):`fiberx.Otel` → `fiberx.Trace` → `fiberx.Recover` → `fiberx.AccessLog` → 业务 handler。`fiberx.Otel` 在 telemetry 未初始化时为 no-op。 `otelx.PrometheusRoutes` 仅在 `metrics.prometheus.enabled: true` 时注册 `/metrics` 端点。`GFX_OTEL_TEST=1` 仅供框架内部测试,业务代码请勿依赖。 本地开发可在 docker-compose 中启动 OpenTelemetry Collector,将 `otlp.endpoint` 指向 `localhost:4317`。 ### Client Instrumentation `postgresx`、`redisx`、`natsx` 和 `rpx` client 支持包级可选 OpenTelemetry tracing 与低基数 metrics。必须同时启用全局 provider 和对应包的 `observability.enabled`; 如果只打开包级开关但 `server.observability.enabled` 未启用,client 正常工作且不采集。 ```yaml server: observability: enabled: true otlp: endpoint: localhost:4317 protocol: grpc insecure: true trace: enabled: true metrics: enabled: true postgres: observability: enabled: true recordDBStatement: false redis: observability: enabled: true nats: observability: enabled: true subjectMode: low_cardinality rpcx: server: observability: enabled: true client: observability: enabled: true ``` 默认不记录 SQL 原文、Redis key、完整 NATS subject、RPC 参数或 payload。 NATS `subjectMode` 支持 `low_cardinality` 和 `none`。Redis 使用 `github.com/redis/go-redis/extra/redisotel/v9@v9.20.0`,与当前 `github.com/redis/go-redis/v9@v9.20.0` 保持同版本线。 ### RPCX Trace Propagation `rpx` can propagate W3C trace context over RPCX request metadata when both sides enable observability and global telemetry is initialized through `otelx`. ```yaml rpc: server: observability: enabled: true client: observability: enabled: true ``` The supported RPC paths are `Client.Call`, `Client.Broadcast`, and `Client.Fork`. `SendRaw`, file transfer, stream, HTTP gateway, and JSON-RPC gateway are outside this propagation layer. ## 服务启动骨架 `server` 包提供基于 Fiber v3 的服务启动骨架。它复用已有的 `config`、 `fiberx`、`result` 和 `trace` 能力,把配置加载、默认中间件、健康检查、依赖 初始化、生命周期日志和优雅关闭收敛到一个 `App[T]` 中。 ### 配置与加载 业务配置可以通过实现 `ServerConfig() server.Config` 暴露服务启动参数。 `server.Load` 会先调用 `config.Load` 读取完整业务配置,再用这份配置创建 静态 Fiber 应用;如果已经在外部拿到了配置,也可以直接调用 `server.New(cfg, ...)`。需要由 App 管理文件监听和模块配置更新时,使用 `server.Watch`。 ```go type AppConfig struct { Server server.Config `yaml:"server"` } func (c AppConfig) ServerConfig() server.Config { return c.Server } configPath := flag.String("config", "configs/config.yaml", "path to config file") flag.Parse() app, err := server.Load[AppConfig]( context.Background(), config.YAMLFile(*configPath), server.WithConfigOptions[AppConfig]( config.WithEnvPrefix[AppConfig]("APP"), ), server.WithServerOptions[AppConfig]( server.WithRoutes[AppConfig](func(app *fiber.App, deps *server.Container) error { app.Get("/hello", func(c fiber.Ctx) error { return fiberx.JSON(c, map[string]string{"message": "hello"}) }) return nil }), ), ) if err != nil { panic(err) } if err := app.Run(context.Background()); err != nil { panic(err) } ``` `WithConfigOptions` 只接收 `config` 包的加载选项,例如环境变量覆盖、额外校验 函数等。`WithServerOptions` 接收服务选项,例如 Fiber 配置、路由、中间件、 初始化器、健康检查器和日志实现。 路由注册器返回 `error`。多个注册器按声明顺序执行;任一个返回错误后,后续注册器 不会执行,服务创建会回滚已初始化的依赖并关闭服务 logger: ```go server.WithRoutes[AppConfig](func(app *fiber.App, deps *server.Container) error { routes, err := buildRoutes(deps) if err != nil { return fmt.Errorf("build routes: %w", err) } registerRoutes(app, routes) return nil }) ``` #### 静态加载与动态监听 - `server.New` 使用调用方已经读取的配置创建静态 App。 - `server.Load` 通过 `config.Load` 读取一次配置,适合配置只在重启时生效的服务。 - `server.Watch` 先完成初始加载和 App 创建,再由 App 持有 watcher;适合需要动态配置的 服务。`Shutdown` 会停止 watcher,并等待正在执行的 Reload 返回后再关闭依赖。 `server.Watch` 接受与 `server.Load` 相同的加载选项: ```go app, err := server.Watch[AppConfig]( context.Background(), config.YAMLFile(*configPath), server.WithConfigOptions[AppConfig]( config.WithEnvPrefix[AppConfig]("APP"), config.WithPollInterval[AppConfig](time.Second), config.WithErrorHandler[AppConfig](func(err error) { log.Printf("config reload failed: %v", err) }), ), server.WithServerOptions[AppConfig](/* initializers and routes */), ) ``` 这里的 `config.WithErrorHandler` 只处理文件 stat、重新读取、环境变量覆盖或配置校验等 watch 错误;模块 `ConfigReloader.Reload` 返回的错误由 `server.App` 写入日志和 readiness 状态,不会传给该 handler。 关闭动态 App 时,等待 Reload 的有效期限取调用方 context deadline 与 `ServerConfig().ShutdownTimeout` 中较早者。`ConfigReloader` 必须观察传入的取消信号并 及时返回;如果它忽略取消直到期限耗尽,`Shutdown` 返回包装后的 deadline error,不会 关闭仍可能被 Reload 使用的 Container 依赖或 Logger。App 保持 closing,`/readyz` 继续返回 `503`,已经停止的 HTTP 服务不会恢复,也不会再启动新 Reload。待回调最终 退出后,可以再次调用 `Shutdown`,继续完成依赖和 Logger 清理。 动态 App 的 `app.Config()` 返回 Store 中通过读取、环境变量覆盖和校验的最新配置 快照;静态 App 的 `Config()` 始终返回启动配置。`app.ServerConfig()` 在两种模式下 都返回启动时解析并应用默认值后的服务配置快照,不会随文件更新。因而即使最新 `Config()` 已包含新的监听地址、HTTP 端口、关闭超时或基础设施 DSN,App 也不会自动 重绑端口、改变当前生命周期参数或重建已有连接。此类字段应通过重启发布,除非业务 为对应模块明确实现了安全的 Reload。 对应 YAML 示例: ```yaml server: name: app host: 0.0.0.0 port: 8080 healthPath: /healthz readinessPath: /readyz shutdownTimeout: 10s ``` ### 默认能力 默认创建的服务会启用: - `fiberx.Trace`:解析或生成 trace id,并写入请求上下文和响应头。 - `fiberx.Recover`:拦截 panic,并按 `result` 统一错误响应返回。 - `GET /healthz`:返回基础存活状态。 - `GET /readyz`:合并应用 closing 状态、配置 Reload 的 failed/skipped 状态和已注册 依赖的健康检查;任一项异常时返回 503。 - 基于 `pkg/logger` 和 `log/slog` 的生命周期日志:记录 initializer 启动/完成/失败、服务监听、 关闭请求、Fiber 关闭、依赖关闭和 logger 关闭等事件。 `server.New` 默认不会启用 `fiberx.AccessLog`,需要业务通过 `WithMiddleware` 显式挂载。传入 `WithoutDefaultMiddleware()` 后,`fiberx.Trace` 和 `fiberx.Recover` 也不会自动挂载,业务需要自行按需注册。 `server.Config` 支持配置服务名、监听地址、健康检查路径、readiness 路径、 关闭超时时间和日志配置。未配置时会使用默认值:服务名 `app`、端口 `8080`、 `/healthz`、`/readyz`、关闭超时 `10s`、日志级别 `info`、文本格式日志、 输出到 `stdout`。 ### 日志配置 `server.log` 通过 `pkg/logger` 创建默认 logger。默认输出到 `stdout`,适合容器 或本地开发: ```yaml server: name: app log: level: info format: text output: stdout ``` 如果需要写入本地文件,可以启用 `output: file`: ```yaml server: name: app log: level: info format: json output: file file: dir: logs rotation: daily ``` 文件日志按日期拆分,路径格式为: ```text logs/YYYY-MM-DD/app.log ``` 第一版 `rotation` 只支持 `daily`。`output` 支持 `stdout`、`stderr` 和 `file`; `format` 支持 `text` 和 `json`;`level` 支持 `debug`、`info`、`warn`、`warning` 和 `error`。 ### 初始化器与依赖生命周期 三方依赖建议通过 `server.Initializer` 或 `server.WithInitializerFunc` 注册。 初始化器会在路由注册前执行,并可以把数据库、缓存、消息队列客户端等对象写入 `server.Container`。 每个 `Initializer` 必须通过 `DependsOn() []string` 声明初始化器名称依赖。框架在 创建 logger、Container 或 Fiber App 前校验 nil、空名称、重复名称、空依赖名称、 缺失依赖、自依赖和循环依赖,并按稳定拓扑顺序初始化:满足依赖关系的前提下, 互不依赖的初始化器仍保持注册顺序。函数初始化器可以用 `WithInitializerFuncDependencies(name, dependencies, fn)` 声明同样的关系。 ```go type Database struct{} func (d *Database) HealthCheck(ctx context.Context) error { return nil } func (d *Database) Close(ctx context.Context) error { return nil } app, err := server.Load[AppConfig]( context.Background(), config.YAMLFile("config.yaml"), server.WithServerOptions[AppConfig]( server.WithInitializerFunc[AppConfig]("database", func(ctx context.Context, cfg AppConfig, deps *server.Container) error { db := &Database{} return deps.Set("database", db) }), server.WithRoutes[AppConfig](func(app *fiber.App, deps *server.Container) error { app.Get("/users", func(c fiber.Ctx) error { db, ok := server.Get[*Database](deps, "database") if !ok { return fiberx.Error(c, result.NewError(fiber.StatusInternalServerError, "database unavailable")) } _ = db return fiberx.JSON(c, []string{}) }) return nil }), ), ) if err != nil { panic(err) } ``` 写入 `Container` 的依赖如果实现了 `HealthCheck(context.Context) error`,会自动 加入 `/readyz` 检查;如果实现了 `Close(context.Context) error`,会在 `Shutdown` 中按注册顺序反向关闭。初始化失败时,已经注册的依赖也会执行清理。 #### 模块配置 Reload 只有初始化器同时实现 `server.ConfigReloader[T]` 时,才会响应 `server.Watch` 产生的 配置更新: ```go func (i *databaseInitializer) Reload( ctx context.Context, previous AppConfig, next AppConfig, deps *server.Container, ) error { // 先完整创建并验证新资源,再以模块自己的同步机制原子替换引用。 // 返回 error 前不能关闭或破坏仍在使用的旧资源。 return nil } ``` Reload 按初始化器的稳定拓扑顺序串行执行,每个模块的 `previous` 是该模块最后一次 成功应用的配置。某模块失败后,依赖它的模块在本轮会跳过;无依赖关系的分支继续 更新。失败或跳过会写入日志,并使 `/readyz` 返回 `503`,但不会终止进程;后续一轮 中所有相关模块成功后 readiness 自动恢复。Store 最新快照可能暂时领先于某个模块 实际生效的配置,这是失败保留旧资源的预期结果。 原子替换由 Reloader 自身负责,框架只负责编排、状态和生命周期。普通 Initializer 不会因为使用了 `server.Watch` 而自动重建;当前内置基础设施初始化器也不应被视为 默认支持 DSN、端口或连接参数热更新。 ## RPCX + Etcd 服务治理基础能力 `etcdx` 管理 etcd v3 client 的配置、健康检查和关闭生命周期;`rpx` 基于 `*etcdx.Client` 创建 rpcx server/client,但不拥有 etcd client 生命周期。RPC 监听端口独立于 HTTP server 端口,rpcx server 默认服务名来自 `server.Config.Name`。 ```yaml server: name: order host: 0.0.0.0 port: 8080 etcd: endpoints: ["127.0.0.1:2379"] dialTimeout: 5s requestTimeout: 3s rpc: server: host: 0.0.0.0 port: 8972 basePath: /rpcx updateInterval: 1m orderClient: servicePath: order basePath: /rpcx inventoryClient: servicePath: inventory basePath: /rpcx ``` ```go type AppConfig struct { Server server.Config `yaml:"server"` Etcd etcdx.Config `yaml:"etcd"` RPC struct { Server rpx.ServerConfig `yaml:"server"` OrderClient rpx.ClientConfig `yaml:"orderClient"` InventoryClient rpx.ClientConfig `yaml:"inventoryClient"` } `yaml:"rpc"` } func (c AppConfig) ServerConfig() server.Config { return c.Server } type OrderService struct{} func (s *OrderService) Create(ctx context.Context, args *CreateOrderArgs, reply *CreateOrderReply) error { return nil } app, err := server.Load[AppConfig]( context.Background(), config.YAMLFile("config.yaml"), server.WithServerOptions[AppConfig]( server.WithInitializer(etcdx.Initializer[AppConfig]("", func(cfg AppConfig) etcdx.Config { return cfg.Etcd })), server.WithInitializer(rpx.ServerInitializer[AppConfig]( "", func(cfg AppConfig) rpx.ServerConfig { return cfg.RPC.Server }, func(cfg AppConfig) server.Config { return cfg.Server }, "", func(ctx context.Context, cfg AppConfig, srv *rpx.Server) error { return srv.Register(new(OrderService)) }, )), server.WithInitializer(rpx.ClientInitializer[AppConfig]( "order_rpcx_client", func(cfg AppConfig) rpx.ClientConfig { return cfg.RPC.OrderClient }, "", )), server.WithInitializer(rpx.ClientInitializer[AppConfig]( "inventory_rpcx_client", func(cfg AppConfig) rpx.ClientConfig { return cfg.RPC.InventoryClient }, "", )), ), ) if err != nil { panic(err) } defer app.Shutdown(context.Background()) ``` 手动创建时,先创建 `etcdx.Client`,再把同一个 client 传给 server/client: ```go etcdClient, err := etcdx.NewClient(cfg.Etcd) if err != nil { panic(err) } defer etcdClient.Close(context.Background()) rpcServer, err := rpx.NewServer( cfg.RPC.Server, etcdClient, rpx.WithDefaultServiceName(cfg.Server.Name), ) if err != nil { panic(err) } if err := rpcServer.Register(new(OrderService)); err != nil { panic(err) } if err := rpcServer.Start(context.Background()); err != nil { panic(err) } defer rpcServer.Close(context.Background()) orderClient, err := rpx.NewClient(cfg.RPC.OrderClient, etcdClient) if err != nil { panic(err) } defer orderClient.Close(context.Background()) var reply CreateOrderReply if err := orderClient.Call(context.Background(), "Create", &CreateOrderArgs{}, &reply); err != nil { panic(err) } ``` 多 client 可以用不同的 container 依赖名注册,例如 `order_rpcx_client` 和 `inventory_rpcx_client`。`rpx.Client.XClient()` 会返回底层 rpcx `client.XClient`, 需要原生高级能力时可以直接使用。