diff --git a/Makefile b/Makefile index 05b9541..f3a8a19 100644 --- a/Makefile +++ b/Makefile @@ -69,11 +69,11 @@ lint: golangci-lint ## Run golangci-lint linter .PHONY: test-database-integration test-database-integration: setup-envtest ## 使用临时 API server、PostgreSQL 与 OpenBao 容器验证 Database 后端。 - KUBEBUILDER_ASSETS="$(shell "$(ENVTEST)" use $(ENVTEST_K8S_VERSION) --bin-dir "$(LOCALBIN)" -p path)" go test -tags=integration -race -count=1 ./internal/database/... + KUBEBUILDER_ASSETS="$(shell "$(ENVTEST)" use $(ENVTEST_K8S_VERSION) --bin-dir "$(LOCALBIN)" -p path)" go test -tags=integration -race -count=1 ./internal/database/... ./internal/infra/... .PHONY: lint-database-integration lint-database-integration: golangci-lint ## 检查集成测试构建标签下的 Database 代码。 - "$(GOLANGCI_LINT)" run --build-tags=integration ./internal/database/... + "$(GOLANGCI_LINT)" run --build-tags=integration ./internal/database/... ./internal/infra/... .PHONY: lint-fix lint-fix: golangci-lint ## Run golangci-lint linter and perform fixes diff --git a/cmd/database.go b/cmd/database.go index 9603c38..75e87fc 100644 --- a/cmd/database.go +++ b/cmd/database.go @@ -9,7 +9,7 @@ import ( ) func setupInstanceObservation(manager ctrl.Manager, namespace, rootCert string) (*application.InstanceService, error) { - credentials, err := kubernetes.NewSecretCredentials(manager.GetConfig(), namespace) + credentials, err := kubernetes.NewSecretCredentials(manager.GetAPIReader(), namespace) if err != nil { return nil, err } diff --git a/cmd/main.go b/cmd/main.go index 4144bd6..6215b1a 100644 --- a/cmd/main.go +++ b/cmd/main.go @@ -42,6 +42,8 @@ func init() { // nolint:gocyclo func main() { + var openBao openBaoOptions + openBao.bindFlags(flag.CommandLine) var databaseNamespace, databaseRootCert string flag.StringVar(&databaseNamespace, "database-secret-namespace", os.Getenv("POD_NAMESPACE"), "固定管理 Secret namespace;为空时不启用 Instance 观测") @@ -179,6 +181,10 @@ func main() { } // +kubebuilder:scaffold:builder + if err := setupOpenBaoAuthentication(mgr, openBao); err != nil { + setupLog.Error(err, "Failed to set up OpenBao authentication") + os.Exit(1) + } var instanceService *application.InstanceService if databaseNamespace != "" { instanceService, err = setupInstanceObservation(mgr, databaseNamespace, databaseRootCert) diff --git a/cmd/openbao.go b/cmd/openbao.go new file mode 100644 index 0000000..6184d24 --- /dev/null +++ b/cmd/openbao.go @@ -0,0 +1,72 @@ +/* +Copyright 2026. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package main + +import ( + "errors" + "flag" + "net/http" + + ctrl "sigs.k8s.io/controller-runtime" + + "git.ddupan.top/panxiao81/ayatori/internal/infra/openbao" +) + +type openBaoOptions struct { + address string + caCert string + mount string + role string + identity openbao.KubernetesIdentity +} + +func (o *openBaoOptions) bindFlags(flags *flag.FlagSet) { + flags.StringVar(&o.address, "openbao-address", "", "OpenBao HTTPS 地址;为空时不启用认证会话") + flags.StringVar(&o.caCert, "openbao-ca-cert", "", "OpenBao 公开 CA PEM 路径;默认使用系统信任根") + flags.StringVar(&o.mount, "openbao-auth-mount", "kubernetes", "OpenBao Kubernetes auth mount") + flags.StringVar(&o.role, "openbao-auth-role", "", "OpenBao 登录 role") + flags.StringVar(&o.identity.Namespace, "openbao-service-account-namespace", "", + "TokenRequest 的固定 ServiceAccount namespace") + flags.StringVar(&o.identity.ServiceAccount, "openbao-service-account-name", "", "TokenRequest 的固定 ServiceAccount 名称") + flags.StringVar(&o.identity.Audience, "openbao-token-audience", "openbao", "SA JWT audience,须匹配 OpenBao role") +} + +func setupOpenBaoAuthentication(manager ctrl.Manager, options openBaoOptions) error { + if options.address == "" { + return nil + } + client, err := openbao.NewClient(options.address, options.caCert) + if err != nil { + return err + } + // 复用 manager 已装配的 Kubernetes client,不重复加载配置或创建客户端。 + session, err := openbao.NewKubernetesSession( + client, manager.GetClient(), options.mount, options.role, options.identity, + ) + if err != nil { + return err + } + if err := manager.Add(session); err != nil { + return err + } + return manager.AddReadyzCheck("openbao-auth", func(_ *http.Request) error { + if !session.Ready() { + return errors.New("OpenBao Kubernetes authentication unavailable") + } + return nil + }) +} diff --git a/config/samples/database_openbao_auth_rbac.yaml b/config/samples/database_openbao_auth_rbac.yaml new file mode 100644 index 0000000..c10e13a --- /dev/null +++ b/config/samples/database_openbao_auth_rbac.yaml @@ -0,0 +1,37 @@ +# 示例:只授权 controller 为固定登录 SA 创建短期 JWT;由管理员替换 namespace/subject 后应用。 +# 不自动纳入 config/default,不包含 kubeconfig、长期 token 或 OpenBao 管理权限。 +apiVersion: v1 +kind: ServiceAccount +metadata: + name: database-openbao-login + namespace: ayatori-system +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: database-openbao-token + namespace: ayatori-system +rules: +- apiGroups: [""] + resources: [serviceaccounts/token] + resourceNames: [database-openbao-login] + verbs: [create] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: database-openbao-token + namespace: ayatori-system +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: database-openbao-token +subjects: +# 集群外示例:对应管理员签发 kubeconfig 的实际用户名,不是登录目标 SA 的名字。 +- kind: User + apiGroup: rbac.authorization.k8s.io + name: ayatori-controller +# 集群内可以改为 manager 自己的 ServiceAccount;不要求两个 subject 同时授权。 +# - kind: ServiceAccount +# name: controller-manager +# namespace: ayatori-system diff --git a/docs/architecture/overview.md b/docs/architecture/overview.md index 5c1485b..7736cef 100644 --- a/docs/architecture/overview.md +++ b/docs/architecture/overview.md @@ -43,6 +43,23 @@ Dev 与 Prod 使用独立的 Kubernetes API、数据库、身份和 controller Proxmox 作为稀缺物理基础设施可以共享,通过 pool、tag、token 和明确的资源范围区分 环境。其他后端尽量使用独立数据库、角色、地址池、DNS 空间与凭据。 +## 进程内依赖边界 + +基础设施能力属于整个 controller-manager,不因首个消费者是 Database 就归入该领域。 +`internal/infra/openbao` 管理官方 SDK client 的 TLS 配置、Kubernetes 认证及 token 生命周期, +不依赖 Database 或其他产品领域。Bao client 默认禁用自动重试,写入结果不确定时由用例处理; +领域适配器不修改共享 client 的全局配置。Kubernetes 客户端、cache 和直连 reader 由 manager 管理; +启动入口负责装配与注入,不在领域适配器内重复创建客户端。 + +读写能力优先直接使用官方 `client.Reader`、`client.Client`、OpenBao KV API 等接口, +不为统一命名再包一层通用 reader/writer,也不引入全局注册中心。共享连接不表示扩大授权; +不同身份或权限边界仍由启动装配显式隔离。 + +领域按用例需要维护 repository 契约,其 adapter 负责 CR/领域对象映射及业务结果转换。 +例如 Database 的七键凭据格式、UID 路径、禁止覆盖和不确定结果处理仍由 Database 维护; +它们不是公共 KV 存储的业务规则。Secret 管理凭据读取注入 `manager.GetAPIReader()`, +保持直连 API server、不缓存 Secret 内容的安全边界;资源写入复用 `manager.GetClient()`。 + ## 数据面 Ayatori 不承载或重新实现数据面。控制面故障只应阻止创建与变更,不应停止已有 VM、 diff --git a/docs/database/README.md b/docs/database/README.md index 468465b..899cc7e 100644 --- a/docs/database/README.md +++ b/docs/database/README.md @@ -168,10 +168,57 @@ Instance 删除首先释放本地连接并撤销 Ready。任何引用它的 Data 真实后端覆盖创建/回读、并发唯一创建、重建适配器读取、软删除冲突、固定前缀 token 拒绝管理路径,以及成功写入后丢失响应;HTTP 故障测试补充不重试和错误脱敏。 -这一切片尚未接入 manager:Kubernetes auth/token 生命周期、Database 状态中的稳定位置和 -已确认步骤、供应 service/controller、PostgreSQL 创建以及 ESO 交付仍未完成。 +认证会话已按下面的显式参数接入 manager;凭据存储尚未接入供应用例。 +Database 状态中的稳定位置和已确认步骤、供应 service/controller、PostgreSQL 创建以及 ESO 交付仍未完成。 测试 token 只用于临时 fixture,不是生产静态 token 配置接口。现有绑定不会触发外部写入。 +## OpenBao Kubernetes 认证会话 + +公共 `internal/infra/openbao.KubernetesSession` 复用官方 Kubernetes auth helper 和 `LifetimeWatcher` +(均为 v2.7.0)。认证直接注入 `manager.GetClient()`,与 reconcile 共用已装配的 Kubernetes +client,不从配置另建客户端。标准 `--kubeconfig` / +`KUBECONFIG` 支持 systemd 或其他集群外运行方式,集群内使用 in-cluster 配置,不要求存在 Pod。 +Kubernetes 身份的签发和更新由部署管理及 client-go 的认证机制负责,不另建 kubeconfig 读取器。 + +每次登录前,通过该 client 的 `SubResource("token").Create` 调用固定 namespace/name 的 +ServiceAccount TokenRequest;写入直连 API server,不读取 cache 或要求额外的 SA get 权限。申请 +audience 匹配 OpenBao role、期望有效期 600 秒的短期 JWT;检查返回值非空且未过期,再交给 +官方 Kubernetes auth helper。JWT 不缓存,不读取投射文件,也不回退静态 OpenBao token; +实际 JWT 有效期由 API server 决定。RBAC 拒绝或 TokenRequest 失败时不会继续 Bao 登录。 +OpenBao 登录结果必须包含有效 token 和有限 TTL。 + +续期、到期阈值与等待时间由 SDK 管理;可续期 token 的续期失败就撤下本地 token,不可续期 +token 由 SDK 监测剩余寿命。会话结束后最多每 5 秒重新登录一次,并重新申请 Kubernetes JWT。 +`Ready()` 仅表示当前 lease 正受 SDK 管理,不授权任何 Database 写入,也不能保证下一次请求 +必然成功。认证/续期响应不写日志、不返回给调用方;后端操作仍独立检查并返回脱敏错误。 + +`Start(ctx)` 退出时清空 client token 并等待续期 goroutine 结束。SDK Stop 不取消已经发出的 +续期 HTTP 请求,因此专用 client 的请求期限固定为 15 秒;不增加新连接池或自己的续期算法。 +同一会话拒绝并发 Start。manager 使用 Runnable 管理生命周期,并增加 `openbao-auth` readiness +检查;认证故障不影响 liveness。`--openbao-address` 为空时不启用,不自动修改生产 auth/RBAC。 +HTTPS 和显式 CA/系统信任根不可通过 BAO 环境变量降级,参数见 [部署合同](deployment.md)。 + +依据官方 [Kubernetes auth](https://openbao.org/docs/auth/kubernetes/) 与 +[token 生命周期](https://openbao.org/docs/concepts/auth/)。真实测试使用 envtest 签发 SA token, +OpenBao 通过专用 reviewer 调用真实 TokenReview;集群外受限 kubeconfig 启动实际 manager, +验证共享 client 申请 JWT、短 TTL 续期、RBAC 撤回/恢复与重新登录、跨 namespace/其他 SA 拒绝、 +错误 OpenBao audience 拒绝以及凭据访问恢复。临时 TokenReview 入口仅允许对应 POST, +两段连接均验证 TLS;其 Docker bridge 入口仅为隔离测试,不修改生产 OpenBao 或 Kubernetes。 +单元测试补充 TokenRequest 失败/空 token 无回退、重新申请 JWT、无期限 lease 拒绝、 +并发生命周期、退出清理及 manager 显式参数不受 BAO 环境身份覆盖。 + +## 公共基础设施与领域适配 + +Bao client/TLS 与认证生命周期已移至 `internal/infra/openbao`,与任何产品领域无关。 +Kubernetes 读写客户端由 manager 管理,Secret 凭据适配器只接收直连的 `client.Reader`。 +`adapter/openbao.Credentials` 仍属于 Database:它直接使用官方 KV v2 API,实现七键凭据、 +UID 路径、CAS=0 与回读确认的领域合同,不把这些规则推广为公共存储语义。 +后续供应用例需要的 repository 接口由领域侧按实际操作定义,不提前增加通用仓储抽象。 +分层约定见[总体架构](../architecture/overview.md#进程内依赖边界)。 + +认证单元测试归公共 infra;真实认证与 Database 凭据读写的组合测试仍在 Database adapter。 +Database 集成测试和 lint 入口同时覆盖 `internal/infra/...`,避免拆包导致 CI 漏测。 + ## 设计入口 - [系统规格](specification.md):规范性行为与验收标准; diff --git a/docs/database/deployment.md b/docs/database/deployment.md index 470dae5..83c1aa8 100644 --- a/docs/database/deployment.md +++ b/docs/database/deployment.md @@ -10,7 +10,7 @@ | 最后更新 | 2026-09-25 | 本文定义 v1alpha1 的运行依赖、启动顺序和部署级配置。Instance 观测已接入 manager; -OpenBao、ESO 与完整供应装配仍是后续实现合同。 +OpenBao 认证可显式启用;ESO 与完整供应装配仍是后续实现合同。 ## 依赖与顺序 @@ -34,18 +34,21 @@ OpenBao、ESO 与完整供应装配仍是后续实现合同。 Instance 观测)与 `--database-root-cert`(公开 PostgreSQL CA PEM 路径)。Deployment 通过 downward API 获取 namespace,Secret 权限由该 namespace 的 Role 授予。 -以下是尚待实现的供应/交付配置合同,不表示当前 manager 接受这些 CLI flags。 +以下表格区分已实现的认证参数与尚待实现的供应/交付参数。 必填项缺失、路径无效或 duration 不为正数时,进程必须在启动 manager 前失败; 不得等到 reconcile 时才逐个资源报告配置错误。 | CLI flag | 必填/默认 | 说明 | | --- | --- | --- | -| `--openbao-address` | 必填 | controller 可访问的 OpenBao API address | +| `--openbao-address` | 已实现,默认空 | HTTPS API 地址;为空时关闭认证会话 | | `--openbao-consumer-address` | 默认同 `--openbao-address` | 写入 Tenant status,必须能被预期外部消费者解析 | -| `--openbao-auth-mount` | `kubernetes` | Kubernetes auth mount 名称 | -| `--openbao-auth-role` | 必填 | controller ServiceAccount 对应 role | +| `--openbao-auth-mount` | 已实现,`kubernetes` | Kubernetes auth mount 名称 | +| `--openbao-auth-role` | 已实现,启用时必填 | OpenBao 登录 role | +| `--openbao-ca-cert` | 已实现,默认系统信任根 | OpenBao 公开 CA PEM 路径 | | `--openbao-kv-mount` | `kv` | KV v2 mount;开发可显式用 `secret` | -| `--openbao-service-account-token-path` | `/var/run/secrets/kubernetes.io/serviceaccount/token` | Kubernetes auth 使用的投射 token 文件 | +| `--openbao-service-account-namespace` | 已实现,启用时必填 | TokenRequest 目标 SA 的固定 namespace | +| `--openbao-service-account-name` | 已实现,启用时必填 | TokenRequest 目标 SA 名称 | +| `--openbao-token-audience` | 已实现,`openbao` | SA JWT audience,必须匹配 OpenBao role | | `--openbao-tenant-base-path` | 默认 `postgresql-tenants` | controller 专属 mount-relative 前缀 | | `--external-secret-store-name` | 必填 | controller 创建的 ExternalSecret 固定引用 | | `--database-root-cert` | 已实现 | 只读 PEM trust bundle,不含私钥;沿用 Instance 连接配置 | @@ -56,6 +59,35 @@ address 必须是绝对 `http` 或 `https` URL,不允许 userinfo、query 或 `/` 开头,不含空段、`.` 或 `..`;base path 还不得编码 KV v2 的 `data`/`metadata` API 层。生产环境的 `--openbao-address` 必须使用 HTTPS;HTTP 只用于明确的开发 fixture。 +### 集群内与 systemd 共用 Kubernetes 认证 + +controller 是 API 客户端,不要求部署为 Pod。OpenBao 认证直接复用 manager 已加载的 +Kubernetes 配置:集群外可用标准 `--kubeconfig`(或 `KUBECONFIG`),集群内可使用 +in-cluster 配置。禁止另要求 `/var/run/secrets/.../token` 文件或解析 kubeconfig 中的 bearer token。 +每次 Bao 登录前通过 TokenRequest 申请新的短期 SA JWT,Kubernetes JWT 与 Bao token 的 +生命周期分别由 API 签发和官方 SDK 续期管理;不把 kubeconfig 本身当成永久有效凭据。 + +管理员为 controller 的实际 Kubernetes 身份授予目标 namespace 内 +`create serviceaccounts/token`,用 `resourceNames` 限定目标 SA;示例见 +[最小 RBAC](../../config/samples/database_openbao_auth_rbac.yaml)。集群外 RoleBinding subject +对应 kubeconfig 的用户/组,集群内可绑定 manager SA;登录目标 SA 可以独立于调用者身份。 +controller 不创建 SA、Role/RoleBinding,也不向自己授予权限。不要求 `get secrets` 来获取 JWT。 +OpenBao role 还需限制 SA 名称、namespace 与 audience,TokenReview reviewer 身份由管理员配置。 + +示例启动参数(仅示意,不包含真实 kubeconfig 或凭据): + +```sh +manager --kubeconfig=/etc/ayatori/controller.kubeconfig \ + --database-secret-namespace=ayatori-system \ + --openbao-address=https://bao.example:8200 \ + --openbao-auth-role=ayatori-database \ + --openbao-service-account-namespace=ayatori-system \ + --openbao-service-account-name=database-openbao-login +``` + +这里的认证成功只开放 manager readiness,不代表 Database 已具备供应或交付能力。 +实例管理 Secret 的 namespace 同样由参数指定,systemd 模式不依赖 `POD_NAMESPACE` 环境变量。 + Tenant 不能选择任意凭据路径。凭据必须能随 Database 保留并安全交付给被授权的新 Tenant; 原 `//` 定位规则不再直接作为新 API 合同。 动态供应位置使用 `/`;导入使用 Database 的显式 credentialRef, diff --git a/docs/database/security.md b/docs/database/security.md index 05ebbe0..ea4b79a 100644 --- a/docs/database/security.md +++ b/docs/database/security.md @@ -24,6 +24,9 @@ Kubernetes 管理员、OpenBao 管理员和 PostgreSQL 管理员是平台信任 ## 凭据处理 - controller 使用 Kubernetes auth 获取短期 OpenBao token,不配置长期静态 token。 +- Kubernetes auth 不等于部署在 Kubernetes 内:复用 manager 的 kubeconfig/in-cluster 身份, + 通过最小 RBAC 的指定 ServiceAccount TokenRequest 获取 JWT,不依赖 Pod 投射文件。 + kubeconfig 的签发、更新与撤销由部署管理负责;申请失败不回退其他机器身份。 - 管理凭据只从 Instance 引用的 controller namespace Secret 读取,不复制到 CR/status/Event/metric/trace;管理员维护 ExternalSecret,由 ESO 同步该 Secret。 - 动态供应密码使用密码学安全随机源;已有可靠关联时复用 OpenBao 现值,结果不确定时停止并报冲突。 diff --git a/go.mod b/go.mod index 3b0931a..dfdec79 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,7 @@ go 1.27.1 require ( github.com/jackc/pgx/v5 v5.11.0 + github.com/openbao/openbao/api/auth/kubernetes/v2 v2.7.0 github.com/openbao/openbao/api/v2 v2.7.0 k8s.io/api v0.37.0 k8s.io/apimachinery v0.37.0 diff --git a/go.sum b/go.sum index 876f46b..7d1e3f1 100644 --- a/go.sum +++ b/go.sum @@ -150,6 +150,8 @@ github.com/onsi/ginkgo/v2 v2.27.4 h1:fcEcQW/A++6aZAZQNUmNjvA9PSOzefMJBerHJ4t8v8Y github.com/onsi/ginkgo/v2 v2.27.4/go.mod h1:ArE1D/XhNXBXCBkKOLkbsb2c81dQHCRcF5zwn/ykDRo= github.com/onsi/gomega v1.39.0 h1:y2ROC3hKFmQZJNFeGAMeHZKkjBL65mIZcvrLQBF9k6Q= github.com/onsi/gomega v1.39.0/go.mod h1:ZCU1pkQcXDO5Sl9/VVEGlDyp+zm0m1cmeG5TOzLgdh4= +github.com/openbao/openbao/api/auth/kubernetes/v2 v2.7.0 h1:Fw/pJRMpMTH83pMByCyikRHhxuBDYcnyiNSiK8OqJW0= +github.com/openbao/openbao/api/auth/kubernetes/v2 v2.7.0/go.mod h1:LkXPq4+8aLyQ+qoNBHcJF7nZFx0PYt2FOu+m7sdpAXU= github.com/openbao/openbao/api/v2 v2.7.0 h1:3CD1l3tr39nQraCgFGAWA5vYvPFzZoZrt3NL7DMQKAc= github.com/openbao/openbao/api/v2 v2.7.0/go.mod h1:uXbMoyH2pjSvNyTepinUvLde8pOJB82EuhUCfOKnKbo= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= diff --git a/internal/database/adapter/kubernetes/credentials.go b/internal/database/adapter/kubernetes/credentials.go index 2e01583..c4df00a 100644 --- a/internal/database/adapter/kubernetes/credentials.go +++ b/internal/database/adapter/kubernetes/credentials.go @@ -22,10 +22,8 @@ import ( "errors" corev1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/validation" - typedcore "k8s.io/client-go/kubernetes/typed/core/v1" - "k8s.io/client-go/rest" + "sigs.k8s.io/controller-runtime/pkg/client" "git.ddupan.top/panxiao81/ayatori/internal/database/application" "git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance" @@ -34,18 +32,16 @@ import ( // SecretCredentials 直接读取 API server,不将 Secret 数据纳入共享 informer cache。 // namespace 在装配时固定,Instance 不能选择跨 namespace 读取。 type SecretCredentials struct { - secrets typedcore.SecretInterface + reader client.Reader + namespace string } -func NewSecretCredentials(config *rest.Config, namespace string) (*SecretCredentials, error) { - if config == nil || len(validation.IsDNS1123Label(namespace)) != 0 { - return nil, errors.New("valid controller namespace and API configuration required") +// NewSecretCredentials 要求注入 manager.GetAPIReader() 或等价直连 reader,不可使用缓存 reader。 +func NewSecretCredentials(reader client.Reader, namespace string) (*SecretCredentials, error) { + if reader == nil || len(validation.IsDNS1123Label(namespace)) != 0 { + return nil, errors.New("valid controller namespace and API reader required") } - client, err := typedcore.NewForConfig(config) - if err != nil { - return nil, application.ErrCredentialsUnavailable - } - return &SecretCredentials{secrets: client.Secrets(namespace)}, nil + return &SecretCredentials{reader: reader, namespace: namespace}, nil } func (r *SecretCredentials) Read(ctx context.Context, ref instance.CredentialReference) (application.Credentials, error) { @@ -53,7 +49,8 @@ func (r *SecretCredentials) Read(ctx context.Context, ref instance.CredentialRef return application.Credentials{}, application.ErrCredentialsInvalid } keys := ref.Values() - secret, err := r.secrets.Get(ctx, keys.Name, metav1.GetOptions{}) + secret := &corev1.Secret{} + err := r.reader.Get(ctx, client.ObjectKey{Namespace: r.namespace, Name: keys.Name}, secret) if err != nil { return application.Credentials{}, application.ErrCredentialsUnavailable } diff --git a/internal/database/adapter/openbao/authentication_integration_test.go b/internal/database/adapter/openbao/authentication_integration_test.go new file mode 100644 index 0000000..8c6401c --- /dev/null +++ b/internal/database/adapter/openbao/authentication_integration_test.go @@ -0,0 +1,361 @@ +//go:build integration + +/* +Copyright 2026. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package openbao_test + +import ( + "context" + "crypto/ecdsa" + "crypto/elliptic" + "crypto/rand" + "crypto/tls" + "crypto/x509" + "encoding/pem" + "math/big" + "net" + "net/http" + "net/http/httptest" + "net/http/httputil" + "net/url" + "os/exec" + "strconv" + "strings" + "testing" + "time" + + bao "github.com/openbao/openbao/api/v2" + authenticationv1 "k8s.io/api/authentication/v1" + corev1 "k8s.io/api/core/v1" + rbacv1 "k8s.io/api/rbac/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/rest" + "k8s.io/client-go/tools/clientcmd" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/envtest" + metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server" + + "git.ddupan.top/panxiao81/ayatori/internal/infra/openbao" +) + +const ( + unrelatedAuthNamespace = "bao-unrelated" + tokenRequestRole = "request-openbao-token" + authNamespace = "bao-controller" + authAudience = "openbao" +) + +// Bao 容器通过 Docker bridge 访问这个仅转发 TokenReview 的临时入口。 +// 上游仍是带 CA 验证的真实 envtest API;不模拟 JWT 签名、audience 或 RBAC 判定。 +func tokenReviewEndpoint(t *testing.T, config *rest.Config) (string, string) { + t.Helper() + upstream, err := url.Parse(config.Host) + if err != nil { + t.Fatal("invalid envtest address") + } + transport, err := rest.TransportFor(rest.AnonymousClientConfig(config)) + if err != nil { + t.Fatal("cannot construct TokenReview transport") + } + proxy := httputil.NewSingleHostReverseProxy(upstream) + proxy.Transport = transport + proxy.ErrorHandler = func(w http.ResponseWriter, _ *http.Request, _ error) { w.WriteHeader(http.StatusBadGateway) } + server := httptest.NewUnstartedServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost || r.URL.Path != "/apis/authentication.k8s.io/v1/tokenreviews" { + w.WriteHeader(http.StatusForbidden) + return + } + proxy.ServeHTTP(w, r) + })) + if err := server.Listener.Close(); err != nil { + t.Fatal("cannot replace fixture listener") + } + server.Listener, err = net.Listen("tcp", "0.0.0.0:0") + if err != nil { + t.Fatal("cannot expose TokenReview fixture") + } + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + output, err := exec.CommandContext(ctx, "docker", "network", "inspect", "bridge", "--format", + `{{(index .IPAM.Config 0).Gateway}}`).Output() + if err != nil { + t.Fatal("cannot locate fixture Docker bridge") + } + gateway := strings.TrimSpace(string(output)) + if net.ParseIP(gateway) == nil { + t.Fatal("invalid fixture bridge gateway") + } + certificate, caPEM := tokenReviewCertificate(t, net.ParseIP(gateway)) + server.TLS = &tls.Config{Certificates: []tls.Certificate{certificate}, MinVersion: tls.VersionTLS12} + server.StartTLS() + t.Cleanup(server.Close) + port := server.Listener.Addr().(*net.TCPAddr).Port + return "https://" + net.JoinHostPort(gateway, strconv.Itoa(port)), caPEM +} + +func tokenReviewCertificate(t *testing.T, address net.IP) (tls.Certificate, string) { + t.Helper() + key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader) + if err != nil { + t.Fatal("cannot create fixture TLS key") + } + template := &x509.Certificate{ + SerialNumber: big.NewInt(1), + NotBefore: time.Now().Add(-time.Minute), + NotAfter: time.Now().Add(time.Hour), + IPAddresses: []net.IP{address}, + KeyUsage: x509.KeyUsageDigitalSignature | x509.KeyUsageCertSign, + ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth}, + IsCA: true, BasicConstraintsValid: true, + } + der, err := x509.CreateCertificate(rand.Reader, template, template, &key.PublicKey, key) + if err != nil { + t.Fatal("cannot create fixture TLS certificate") + } + certificate := tls.Certificate{Certificate: [][]byte{der}, PrivateKey: key} + return certificate, string(pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der})) +} + +type kubernetesAuthFixture struct { + config *rest.Config + controllerConfig *rest.Config + admin *kubernetes.Clientset + controller *kubernetes.Clientset + grant func() + requestToken func(string, []string) string +} + +func newKubernetesAuthFixture(t *testing.T) *kubernetesAuthFixture { + t.Helper() + environment := &envtest.Environment{} + config, err := environment.Start() + if err != nil { + t.Fatal("cannot start authentication API fixture", err) + } + t.Cleanup(func() { + if err := environment.Stop(); err != nil { + t.Error("cannot stop authentication API fixture") + } + }) + clientset, err := kubernetes.NewForConfig(config) + if err != nil { + t.Fatal("cannot construct fixture API client") + } + ctx := t.Context() + for _, namespace := range []string{authNamespace, unrelatedAuthNamespace} { + if _, err := clientset.CoreV1().Namespaces().Create(ctx, &corev1.Namespace{Name: namespace}, metav1.CreateOptions{}); err != nil { + t.Fatal("cannot create fixture namespace") + } + if _, err := clientset.CoreV1().ServiceAccounts(namespace).Create(ctx, &corev1.ServiceAccount{Name: authRole}, metav1.CreateOptions{}); err != nil { + t.Fatal("cannot create fixture service account") + } + } + // 集群外身份只有指定 SA 的 TokenRequest 权限,不使用 envtest 管理员凭据运行会话。 + user, err := environment.AddUser(envtest.User{Name: "systemd-controller"}, config) + if err != nil { + t.Fatal("cannot create external controller identity") + } + kubeconfig, err := user.KubeConfig() + if err != nil { + t.Fatal("cannot build external kubeconfig") + } + externalConfig, err := clientcmd.RESTConfigFromKubeConfig(kubeconfig) + if err != nil { + t.Fatal("cannot load external kubeconfig") + } + if _, err := clientset.RbacV1().Roles(authNamespace).Create(ctx, &rbacv1.Role{ + Name: tokenRequestRole, + Rules: []rbacv1.PolicyRule{{APIGroups: []string{""}, Resources: []string{"serviceaccounts/token"}, + ResourceNames: []string{authRole}, Verbs: []string{"create"}}}, + }, metav1.CreateOptions{}); err != nil { + t.Fatal("cannot create TokenRequest Role") + } + permission := &rbacv1.RoleBinding{ + Name: tokenRequestRole, + RoleRef: rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "Role", Name: tokenRequestRole}, + Subjects: []rbacv1.Subject{{Kind: "User", APIGroup: rbacv1.GroupName, Name: "systemd-controller"}}, + } + grant := func() { + t.Helper() + if _, err := clientset.RbacV1().RoleBindings(authNamespace).Create(ctx, permission.DeepCopy(), metav1.CreateOptions{}); err != nil { + t.Fatal("cannot grant TokenRequest permission") + } + } + grant() + externalClient, err := kubernetes.NewForConfig(externalConfig) + if err != nil { + t.Fatal("cannot construct restricted controller client") + } + for _, target := range []struct{ namespace, name string }{{unrelatedAuthNamespace, authRole}, {authNamespace, "another-account"}} { + _, err := externalClient.CoreV1().ServiceAccounts(target.namespace).CreateToken(ctx, target.name, + &authenticationv1.TokenRequest{Spec: authenticationv1.TokenRequestSpec{Audiences: []string{authAudience}}}, metav1.CreateOptions{}) + if !apierrors.IsForbidden(err) { + t.Fatal("TokenRequest escaped Role namespace/resourceNames restriction") + } + } + if _, err := clientset.RbacV1().ClusterRoleBindings().Create(ctx, &rbacv1.ClusterRoleBinding{ + Name: "bao-fixture-reviewer", + RoleRef: rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "ClusterRole", Name: "system:auth-delegator"}, + Subjects: []rbacv1.Subject{{Kind: "ServiceAccount", Namespace: authNamespace, Name: authRole}}, + }, metav1.CreateOptions{}); err != nil { + t.Fatal("cannot authorize fixture TokenReview") + } + requestToken := func(namespace string, audiences []string) string { + t.Helper() + response, err := clientset.CoreV1().ServiceAccounts(namespace).CreateToken(ctx, authRole, + &authenticationv1.TokenRequest{Spec: authenticationv1.TokenRequestSpec{Audiences: audiences}}, metav1.CreateOptions{}) + if err != nil { + t.Fatal("cannot issue fixture service account token") + } + return response.Status.Token + } + + return &kubernetesAuthFixture{ + config: config, controllerConfig: externalConfig, admin: clientset, controller: externalClient, + grant: grant, requestToken: requestToken, + } +} + +const authRole = "controller" + +var testIdentity = openbao.KubernetesIdentity{Namespace: authNamespace, ServiceAccount: authRole, Audience: authAudience} + +func waitForAuthentication(t *testing.T, check func() bool) { + t.Helper() + deadline := time.NewTimer(20 * time.Second) + defer deadline.Stop() + for !check() { + select { + case <-deadline.C: + t.Fatal("authentication condition timed out") + case <-time.After(20 * time.Millisecond): + } + } +} + +func startSession(t *testing.T, session interface{ Start(context.Context) error }) context.CancelFunc { + t.Helper() + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan error, 1) + go func() { done <- session.Start(ctx) }() + t.Cleanup(func() { + cancel() + select { + case <-done: + case <-time.After(20 * time.Second): + t.Error("authentication did not stop") + } + }) + return cancel +} + +// 此处验证公共认证会话与 Database 凭据适配器的跨层集成。 +func TestKubernetesSessionWithRealTokenReview(t *testing.T) { + api := newKubernetesAuthFixture(t) + ctx := t.Context() + root := baoFixture(t) + if err := root.Sys().EnableAuthWithOptionsWithContext(ctx, "kubernetes", &bao.EnableAuthOptions{Type: "kubernetes"}); err != nil { + t.Fatal("cannot enable fixture Kubernetes auth") + } + reviewerToken := api.requestToken(authNamespace, nil) + reviewURL, reviewCA := tokenReviewEndpoint(t, api.config) + if _, err := root.Logical().WriteWithContext(ctx, "auth/kubernetes/config", map[string]any{ + "kubernetes_host": reviewURL, + "kubernetes_ca_cert": reviewCA, + "token_reviewer_jwt": reviewerToken, + "disable_local_ca_jwt": true, + }); err != nil { + t.Fatal("cannot configure fixture TokenReview:", strings.NewReplacer(reviewerToken, "[REDACTED]", fixtureToken, "[REDACTED]").Replace(err.Error())) + } + if err := root.Sys().PutPolicyWithContext(ctx, authRole, `path "secret/data/applications/*" { capabilities = ["create", "update", "read"] }`); err != nil { + t.Fatal("cannot configure fixture credential policy") + } + if _, err := root.Logical().WriteWithContext(ctx, "auth/kubernetes/role/controller", map[string]any{ + "bound_service_account_names": []string{authRole}, + "bound_service_account_namespaces": []string{authNamespace}, + "audience": authAudience, + "token_policies": []string{authRole}, + "token_ttl": "3s", "token_max_ttl": "60s", + }); err != nil { + t.Fatal("cannot configure fixture auth role") + } + client := fixtureClient(t, root.Address()) + manager, err := ctrl.NewManager(api.controllerConfig, ctrl.Options{Metrics: metricsserver.Options{BindAddress: "0"}, HealthProbeBindAddress: "0"}) + if err != nil { + t.Fatal("cannot construct external controller manager") + } + session, err := openbao.NewKubernetesSession(client, manager.GetClient(), "kubernetes", authRole, testIdentity) + if err != nil { + t.Fatal(err) + } + if err := manager.Add(session); err != nil { + t.Fatal("cannot register authentication lifecycle") + } + cancel := startSession(t, manager) + waitForAuthentication(t, session.Ready) + store := fixtureStore(t, client) + if err := store.Create(ctx, credentialPath, fixtureCredential(t)); err != nil { + t.Fatal("Kubernetes identity cannot create scoped credential", err) + } + initialToken := client.Token() + // 超过初始 TTL 后同一 token 仍可用,证明发生真实 renew-self,而非只登录一次。 + start := time.Now() + waitForAuthentication(t, func() bool { return time.Since(start) > 4*time.Second }) + if client.Token() != initialToken { + t.Fatal("token was replaced before renewal could be verified") + } + if _, err := client.Auth().Token().LookupSelfWithContext(ctx); err != nil { + t.Fatal("short-lived token was not renewed") + } + for _, invalid := range []struct { + namespace string + audience string + }{ + {unrelatedAuthNamespace, authAudience}, {authNamespace, "wrong-audience"}, + } { + // 用相同真实登录入口直接确认 namespace/audience 拒绝,不依赖定时轮询推断。 + if _, err := root.Logical().WriteWithContext(ctx, "auth/kubernetes/login", map[string]any{ + "role": authRole, "jwt": api.requestToken(invalid.namespace, []string{invalid.audience}), + }); err == nil { + t.Fatal("invalid Kubernetes identity was accepted") + } + } + if err := api.admin.RbacV1().RoleBindings(authNamespace).Delete(ctx, tokenRequestRole, metav1.DeleteOptions{}); err != nil { + t.Fatal("cannot revoke TokenRequest permission") + } + if err := root.Auth().Token().RevokeOrphanWithContext(ctx, client.Token()); err != nil { + t.Fatal("cannot revoke fixture OpenBao token") + } + waitForAuthentication(t, func() bool { return !session.Ready() && client.Token() == "" }) + denied, err := api.controller.CoreV1().ServiceAccounts(authNamespace).CreateToken(ctx, authRole, + &authenticationv1.TokenRequest{Spec: authenticationv1.TokenRequestSpec{Audiences: []string{authAudience}}}, metav1.CreateOptions{}) + if !apierrors.IsForbidden(err) || (denied != nil && denied.Status.Token != "") { + t.Fatal("revoked caller still obtained a token") + } + api.grant() + waitForAuthentication(t, session.Ready) + if client.Token() == initialToken { + t.Fatal("reauthentication reused revoked OpenBao token") + } + if _, err := store.Read(ctx, credentialPath); err != nil { + t.Fatal("credential access did not recover", err) + } + cancel() + waitForAuthentication(t, func() bool { return !session.Ready() && client.Token() == "" }) +} diff --git a/internal/database/adapter/openbao/credentials.go b/internal/database/adapter/openbao/credentials.go index ff80459..ecfddd8 100644 --- a/internal/database/adapter/openbao/credentials.go +++ b/internal/database/adapter/openbao/credentials.go @@ -49,12 +49,11 @@ type Credentials struct { } // NewCredentials 不登录、不读取环境 token。调用方必须提供专用的已认证 client。 -// 禁用 SDK 写入重试,防止第一次结果丢失后被 CAS 错误掩盖。 +// client 由公共 infra 禁用自动重试,防止第一次结果丢失后被 CAS 错误掩盖。 func NewCredentials(client *bao.Client, mount, basePath string) (*Credentials, error) { if client == nil || !validPath(mount) || !validPath(basePath) { return nil, ErrInvalidLocation } - client.SetMaxRetries(0) return &Credentials{kv: client.KVv2(mount), basePath: basePath}, nil } diff --git a/internal/database/adapter/openbao/credentials_test.go b/internal/database/adapter/openbao/credentials_test.go index 26cce76..b164173 100644 --- a/internal/database/adapter/openbao/credentials_test.go +++ b/internal/database/adapter/openbao/credentials_test.go @@ -104,6 +104,7 @@ func fixtureClient(t *testing.T, address string) *bao.Client { t.Helper() config := bao.DefaultConfig() config.Address = address + config.MaxRetries = 0 client, err := bao.NewClient(config) if err != nil { t.Fatal("cannot construct fixture client") diff --git a/internal/database/adapter/postgresql/fixture_integration_test.go b/internal/database/adapter/postgresql/fixture_integration_test.go index 593c487..9e37531 100644 --- a/internal/database/adapter/postgresql/fixture_integration_test.go +++ b/internal/database/adapter/postgresql/fixture_integration_test.go @@ -32,6 +32,7 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" + kubeclient "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/envtest" secretadapter "git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes" @@ -244,7 +245,11 @@ func newCredentialFixture(t *testing.T) *credentialFixture { } } - reader, err := secretadapter.NewSecretCredentials(config, controllerNamespace) + apiReader, err := kubeclient.New(config, kubeclient.Options{}) + if err != nil { + t.Fatal(err) + } + reader, err := secretadapter.NewSecretCredentials(apiReader, controllerNamespace) if err != nil { t.Fatal(err) } @@ -252,7 +257,11 @@ func newCredentialFixture(t *testing.T) *credentialFixture { if err != nil { t.Fatal(err) } - deniedReader, err := secretadapter.NewSecretCredentials(user.Config(), controllerNamespace) + deniedAPIReader, err := kubeclient.New(user.Config(), kubeclient.Options{}) + if err != nil { + t.Fatal(err) + } + deniedReader, err := secretadapter.NewSecretCredentials(deniedAPIReader, controllerNamespace) if err != nil { t.Fatal(err) } diff --git a/internal/database/adapter/postgresql/instance_controller_integration_test.go b/internal/database/adapter/postgresql/instance_controller_integration_test.go index 78e8f5c..c6e9d57 100644 --- a/internal/database/adapter/postgresql/instance_controller_integration_test.go +++ b/internal/database/adapter/postgresql/instance_controller_integration_test.go @@ -50,14 +50,6 @@ func TestInstanceControllerWithRealPostgreSQL(t *testing.T) { t.Fatal(err) } restricted := instanceControllerRBAC(t, f, apiClient) - credentials, err := secretadapter.NewSecretCredentials(restricted, controllerNamespace) - if err != nil { - t.Fatal(err) - } - service, err := application.NewInstanceService(credentials, postgresql.Connector{}) - if err != nil { - t.Fatal(err) - } // 同进程 -count 重复启动测试 manager;生产继续校验 controller 名称唯一。 skipRepeatedName := true manager, err := ctrl.NewManager(restricted, ctrl.Options{ @@ -68,6 +60,14 @@ func TestInstanceControllerWithRealPostgreSQL(t *testing.T) { if err != nil { t.Fatal(err) } + credentials, err := secretadapter.NewSecretCredentials(manager.GetAPIReader(), controllerNamespace) + if err != nil { + t.Fatal(err) + } + service, err := application.NewInstanceService(credentials, postgresql.Connector{}) + if err != nil { + t.Fatal(err) + } reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: controllerNamespace} if err := reconciler.SetupWithManager(manager); err != nil { t.Fatal(err) diff --git a/internal/infra/openbao/authentication.go b/internal/infra/openbao/authentication.go new file mode 100644 index 0000000..9918ebd --- /dev/null +++ b/internal/infra/openbao/authentication.go @@ -0,0 +1,192 @@ +/* +Copyright 2026. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package openbao + +import ( + "context" + "errors" + "regexp" + "strings" + "sync" + "sync/atomic" + "time" + + kubernetesauth "github.com/openbao/openbao/api/auth/kubernetes/v2" + bao "github.com/openbao/openbao/api/v2" + authenticationv1 "k8s.io/api/authentication/v1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/util/validation" + kubeclient "sigs.k8s.io/controller-runtime/pkg/client" +) + +var ErrAuthenticationConfiguration = errors.New("invalid OpenBao Kubernetes authentication configuration") + +var authPathSegment = regexp.MustCompile(`^[A-Za-z0-9_-]+$`) + +func validAuthMount(mount string) bool { + for segment := range strings.SplitSeq(mount, "/") { + if !authPathSegment.MatchString(segment) { + return false + } + } + return true +} + +// KubernetesSession 为专用 SDK client 维护短期登录,不持久化或对外返回 token。 +// 每次登录通过 manager 的 Kubernetes 身份申请新的 SA JWT,不依赖 controller 的部署位置。 +// 续期调度由官方 LifetimeWatcher 负责,不实现自己的 lease 算法。 +type KubernetesSession struct { + client *bao.Client + mount string + role string + kubernetes kubeclient.Client + identity KubernetesIdentity + running sync.Mutex + ready atomic.Bool +} + +// KubernetesIdentity 是部署固定的登录目标,不由业务请求选择。 +type KubernetesIdentity struct { + Namespace string + ServiceAccount string + Audience string +} + +func NewKubernetesSession( + client *bao.Client, + kubernetes kubeclient.Client, + mount, role string, + identity KubernetesIdentity, +) (*KubernetesSession, error) { + if client == nil || kubernetes == nil || !validAuthMount(mount) || !authPathSegment.MatchString(role) || + len(validation.IsDNS1123Label(identity.Namespace)) != 0 || + len(validation.IsDNS1123Subdomain(identity.ServiceAccount)) != 0 || strings.TrimSpace(identity.Audience) == "" { + return nil, ErrAuthenticationConfiguration + } + client.ClearToken() + client.SetMaxRetries(0) + client.SetClientTimeout(15 * time.Second) + return &KubernetesSession{ + client: client, + mount: mount, + role: role, + kubernetes: kubernetes, + identity: identity, + }, nil +} + +// Ready 仅代表当前登录 lease 正由 SDK 管理,不保证下一次后端请求一定成功。 +func (s *KubernetesSession) Ready() bool { return s.ready.Load() } + +// Start 可交给 manager 管理;关闭时清空本地 token,不撤销共享后端数据。 +// SDK 的 Stop 不取消已发出的续期 HTTP 请求,因此等待该请求结束后才退出,最长受 client timeout 限制。 +func (s *KubernetesSession) Start(ctx context.Context) error { + if !s.running.TryLock() { + return errors.New("OpenBao authentication is already running") + } + defer s.running.Unlock() + defer s.clear() + for ctx.Err() == nil { + s.clear() + secret := s.login(ctx) + if secret != nil { + s.watch(ctx, secret) + } + s.clear() + // 登录失败及不能继续续期均有限速,防止依赖故障时形成请求忙循环。 + select { + case <-ctx.Done(): + return nil + case <-time.After(5 * time.Second): + } + } + return nil +} + +func (s *KubernetesSession) clear() { + s.ready.Store(false) + s.client.ClearToken() +} + +func (s *KubernetesSession) login(ctx context.Context) *bao.Secret { + ctx, cancel := context.WithTimeout(ctx, 15*time.Second) + defer cancel() + // JWT 只用于本次登录,不缓存或自行解析 kubeconfig 中的凭据。 + // client-go 负责 kubeconfig/in-cluster 身份与凭据更新;API server 按 RBAC 签发。 + expirationSeconds := int64(600) + account := &corev1.ServiceAccount{ + Namespace: s.identity.Namespace, + Name: s.identity.ServiceAccount, + } + token := &authenticationv1.TokenRequest{ + Spec: authenticationv1.TokenRequestSpec{ + Audiences: []string{s.identity.Audience}, + ExpirationSeconds: &expirationSeconds, + }, + } + // 子资源写入直接请求 API server,不读 cache,也不需要额外的 ServiceAccount get 权限。 + err := s.kubernetes.SubResource("token").Create(ctx, account, token) + if err != nil || strings.TrimSpace(token.Status.Token) == "" || + !token.Status.ExpirationTimestamp.After(time.Now()) { + return nil + } + // helper 会缓存 token,不能跨登录轮次复用。 + method, err := kubernetesauth.NewKubernetesAuth(s.role, + kubernetesauth.WithMountPath(s.mount), + kubernetesauth.WithServiceAccountToken(token.Status.Token), + ) + if err != nil { + return nil + } + // 先检查短期 lease,再发布到共享 client,避免暴露不合规的登录结果。 + secret, err := method.Login(ctx, s.client) + if err != nil || secret == nil || secret.Auth == nil || secret.Auth.ClientToken == "" || secret.Auth.LeaseDuration <= 0 { + return nil + } + if ctx.Err() != nil { + return nil + } + s.client.SetToken(secret.Auth.ClientToken) + return secret +} + +func (s *KubernetesSession) watch(ctx context.Context, secret *bao.Secret) { + behavior := bao.RenewBehaviorErrorOnErrors + if !secret.Auth.Renewable { + behavior = bao.RenewBehaviorRenewDisabled + } + watcher, err := s.client.NewLifetimeWatcher(&bao.LifetimeWatcherInput{Secret: secret, RenewBehavior: behavior}) + if err != nil { + return + } + s.ready.Store(true) + go watcher.Start() + for { + select { + case <-ctx.Done(): + s.clear() + watcher.Stop() + <-watcher.DoneCh() + return + case <-watcher.DoneCh(): + watcher.Stop() + return + case <-watcher.RenewCh(): + // 不记录 SDK Secret 或 token;无需复制 SDK 已处理的续期数据。 + } + } +} diff --git a/internal/infra/openbao/authentication_test.go b/internal/infra/openbao/authentication_test.go new file mode 100644 index 0000000..779cfb8 --- /dev/null +++ b/internal/infra/openbao/authentication_test.go @@ -0,0 +1,233 @@ +/* +Copyright 2026. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package openbao_test + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" + + bao "github.com/openbao/openbao/api/v2" + authenticationv1 "k8s.io/api/authentication/v1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/rest" + kubeclient "sigs.k8s.io/controller-runtime/pkg/client" + + "git.ddupan.top/panxiao81/ayatori/internal/infra/openbao" +) + +const ( + authRole = "controller" + fixtureToken = "AYATORI-TEST-ONLY-bao-token" +) + +func fixtureClient(t *testing.T, address string) *bao.Client { + t.Helper() + config := bao.NewConfig() + config.Address = address + client, err := bao.NewClient(config) + if err != nil { + t.Fatal("cannot construct fixture client") + } + client.SetToken(fixtureToken) + return client +} + +var testIdentity = openbao.KubernetesIdentity{Namespace: "bao-controller", ServiceAccount: authRole, Audience: "openbao"} + +func authenticationClient(t *testing.T, address string) kubeclient.Client { + t.Helper() + // HTTP 单元 fixture 只提供 TokenRequest;静态映射避免额外模拟 discovery API。 + mapper := meta.NewDefaultRESTMapper([]schema.GroupVersion{corev1.SchemeGroupVersion}) + mapper.Add(corev1.SchemeGroupVersion.WithKind("ServiceAccount"), meta.RESTScopeNamespace) + client, err := kubeclient.New(&rest.Config{Host: address, ContentType: "application/json"}, kubeclient.Options{ + Mapper: mapper, + }) + if err != nil { + t.Fatal("cannot construct fixture Kubernetes client") + } + return client +} + +func authenticationServer(t *testing.T, token func() (string, bool), login http.HandlerFunc) *httptest.Server { + t.Helper() + return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/api/v1/namespaces/bao-controller/serviceaccounts/controller/token" { + login(w, r) + return + } + var request authenticationv1.TokenRequest + if json.NewDecoder(r.Body).Decode(&request) != nil || r.Method != http.MethodPost || + len(request.Spec.Audiences) != 1 || request.Spec.Audiences[0] != testIdentity.Audience || + request.Spec.ExpirationSeconds == nil || *request.Spec.ExpirationSeconds != 600 { + t.Error("unexpected TokenRequest target or lifetime") + w.WriteHeader(http.StatusBadRequest) + return + } + jwt, allowed := token() + if !allowed { + w.WriteHeader(http.StatusForbidden) + return + } + w.Header().Set("Content-Type", "application/json") + if err := json.NewEncoder(w).Encode(authenticationv1.TokenRequest{ + APIVersion: "authentication.k8s.io/v1", Kind: "TokenRequest", + Status: authenticationv1.TokenRequestStatus{Token: jwt, ExpirationTimestamp: metav1.NewTime(time.Now().Add(10 * time.Minute))}, + }); err != nil { + t.Error("cannot encode fixture TokenRequest response") + } + })) +} + +func TestKubernetesSessionDoesNotFallbackFromMissingToken(t *testing.T) { + var requests atomic.Int32 + for _, missing := range []bool{true, false} { + server := authenticationServer(t, func() (string, bool) { return "", !missing }, func(w http.ResponseWriter, _ *http.Request) { + requests.Add(1) + w.WriteHeader(http.StatusForbidden) + }) + defer server.Close() + client := fixtureClient(t, server.URL) + session, err := openbao.NewKubernetesSession(client, authenticationClient(t, server.URL), "kubernetes", authRole, testIdentity) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond) + err = session.Start(ctx) + cancel() + if err != nil || session.Ready() || client.Token() != "" || requests.Load() != 0 { + t.Fatal("denied or empty TokenRequest must not fall back to another identity") + } + } +} + +func waitForAuthentication(t *testing.T, check func() bool) { + t.Helper() + deadline := time.NewTimer(20 * time.Second) + defer deadline.Stop() + for !check() { + select { + case <-deadline.C: + t.Fatal("authentication condition timed out") + case <-time.After(20 * time.Millisecond): + } + } +} + +func startSession(t *testing.T, session interface{ Start(context.Context) error }) context.CancelFunc { + t.Helper() + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan error, 1) + go func() { done <- session.Start(ctx) }() + t.Cleanup(func() { + cancel() + select { + case <-done: + case <-time.After(20 * time.Second): + t.Error("authentication did not stop") + } + }) + return cancel +} + +func TestKubernetesSessionRequestsNewTokenAfterFailure(t *testing.T) { + var attempts atomic.Int32 + var accepted atomic.Bool + var rotated atomic.Bool + server := authenticationServer(t, func() (string, bool) { + if rotated.Load() { + return "rotated-test-jwt", true + } + return "expired-test-jwt", true + }, func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/v1/auth/kubernetes/login" { + w.WriteHeader(http.StatusNotFound) + return + } + var request struct{ JWT, Role string } + if json.NewDecoder(r.Body).Decode(&request) != nil || request.Role != authRole { + t.Error("unexpected authentication request") + w.WriteHeader(http.StatusBadRequest) + return + } + attempts.Add(1) + if request.JWT != "rotated-test-jwt" { + w.WriteHeader(http.StatusForbidden) + return + } + accepted.Store(true) + if err := json.NewEncoder(w).Encode(map[string]any{"auth": map[string]any{ + "client_token": fixtureToken, "lease_duration": 60, "renewable": false, + }}); err != nil { + t.Error("cannot encode authentication fixture response") + } + }) + defer server.Close() + client := fixtureClient(t, server.URL) + session, err := openbao.NewKubernetesSession(client, authenticationClient(t, server.URL), "kubernetes", authRole, testIdentity) + if err != nil { + t.Fatal(err) + } + if client.Token() != "" { + t.Fatal("constructor retained preexisting static token") + } + cancel := startSession(t, session) + waitForAuthentication(t, func() bool { return attempts.Load() > 0 }) + if session.Ready() || client.Token() != "" { + t.Fatal("failed login retained credentials") + } + rotated.Store(true) + waitForAuthentication(t, func() bool { return accepted.Load() && session.Ready() }) + if client.Token() != fixtureToken { + t.Fatal("successful login did not configure the client") + } + if session.Start(t.Context()) == nil { + t.Fatal("allowed concurrent lifecycle owners") + } + cancel() + waitForAuthentication(t, func() bool { return !session.Ready() && client.Token() == "" }) +} + +func TestKubernetesSessionRejectsUnboundedLease(t *testing.T) { + var attempts atomic.Int32 + server := authenticationServer(t, func() (string, bool) { return "test-jwt", true }, func(w http.ResponseWriter, _ *http.Request) { + attempts.Add(1) + if err := json.NewEncoder(w).Encode(map[string]any{"auth": map[string]any{ + "client_token": fixtureToken, "lease_duration": 0, + }}); err != nil { + t.Error("cannot encode authentication fixture response") + } + }) + defer server.Close() + client := fixtureClient(t, server.URL) + session, err := openbao.NewKubernetesSession(client, authenticationClient(t, server.URL), "kubernetes", authRole, testIdentity) + if err != nil { + t.Fatal(err) + } + startSession(t, session) + waitForAuthentication(t, func() bool { return attempts.Load() > 0 }) + if session.Ready() || client.Token() != "" { + t.Fatal("accepted a token without a finite lease") + } +} diff --git a/internal/infra/openbao/client.go b/internal/infra/openbao/client.go new file mode 100644 index 0000000..70b7715 --- /dev/null +++ b/internal/infra/openbao/client.go @@ -0,0 +1,51 @@ +/* +Copyright 2026. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Package openbao 管理控制面共享的 OpenBao 连接与认证,不依赖任何产品领域。 +package openbao + +import ( + "errors" + "net/url" + "strings" + "time" + + bao "github.com/openbao/openbao/api/v2" +) + +// NewClient 创建供读写适配器共享的官方 SDK client;身份由 KubernetesSession 管理。 +func NewClient(addressURL, caCert string) (*bao.Client, error) { + address, err := url.Parse(addressURL) + if err != nil || address.Scheme != "https" || address.Host == "" || address.User != nil || + address.RawQuery != "" || address.ForceQuery || address.Fragment != "" || + (address.Path != "" && address.Path != "/") { + return nil, errors.New("OpenBao address must be an absolute HTTPS URL without credentials, query, fragment or path") + } + // NewConfig 不读取 BAO_TOKEN/BAO_SKIP_VERIFY 等环境配置,不允许旁路 Kubernetes 身份或 TLS。 + config := bao.NewConfig() + config.Address = strings.TrimSuffix(addressURL, "/") + // 公共读写 client 不自动重试:写入结果不确定时交由具体用例决定恢复行为。 + config.MaxRetries = 0 + config.Timeout = 15 * time.Second + if config.Error != nil || config.ConfigureTLS(&bao.TLSConfig{CACert: caCert}) != nil { + return nil, errors.New("cannot configure OpenBao TLS trust") + } + client, err := bao.NewClient(config) + if err != nil { + return nil, errors.New("cannot construct OpenBao client") + } + return client, nil +} diff --git a/internal/infra/openbao/client_test.go b/internal/infra/openbao/client_test.go new file mode 100644 index 0000000..fa15381 --- /dev/null +++ b/internal/infra/openbao/client_test.go @@ -0,0 +1,67 @@ +/* +Copyright 2026. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package openbao_test + +import ( + "net/http" + "net/http/httptest" + "testing" + + "git.ddupan.top/panxiao81/ayatori/internal/infra/openbao" +) + +func TestOpenBaoClientDoesNotUseEnvironmentIdentityOrAddress(t *testing.T) { + t.Setenv("BAO_TOKEN", "TEST-ONLY-unwanted-static-token") + t.Setenv("BAO_ADDR", "http://unwanted.invalid") + t.Setenv("BAO_SKIP_VERIFY", "true") + client, err := openbao.NewClient("https://bao.example/", "") + if err != nil { + t.Fatal(err) + } + if client.Address() != "https://bao.example" || client.Token() != "" { + t.Fatal("ambient environment replaced the explicit connection or identity") + } + if client.MaxRetries() != 0 { + t.Fatal("shared client must not automatically retry uncertain writes") + } + server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + })) + defer server.Close() + untrusted, err := openbao.NewClient(server.URL, "") + if err != nil { + t.Fatal(err) + } + untrusted.SetMaxRetries(0) + if _, err := untrusted.Sys().HealthWithContext(t.Context()); err == nil { + t.Fatal("BAO_SKIP_VERIFY bypassed TLS validation") + } +} + +func TestOpenBaoClientRejectsUnsafeConfiguration(t *testing.T) { + for _, address := range []string{ + "", "http://bao.example", "https://user:secret@bao.example", + "https://bao.example/?token=secret", "https://bao.example/#secret", "https://bao.example/path", + } { + if _, err := openbao.NewClient(address, ""); err == nil { + t.Fatal("accepted unsafe OpenBao address") + } + } + if _, err := openbao.NewClient("https://bao.example", "/nonexistent/fixture-ca"); err == nil { + t.Fatal("accepted missing explicit CA") + } +}