已合并
feat(service):上下文能力demo、集成Kubernetes至上下文 #418
m0u55e创建于 8月14日
feat(service):上下文能力demo、集成Kubernetes至上下文 #418
已合并
共 20 个文件变更+2321-3
| @@ -0,0 +1,264 @@ | |||
| 1 | +# Service 简单能力 Demo | ||
| 2 | + | ||
| 3 | +该应用通过八个 HTTP 接口展示数据库读写、Redis 读写、Envelope 解析和 Kubernetes Pod 生命周期 | ||
| 4 | +能力。本地配置使用 SQLite、FakeRedis 和进程内 Fake Kubernetes。服务器配置使用 MySQL、Redis 和 | ||
| 5 | +`kubernetes_asyncio`。两种配置共用 Handler、请求契约和接口路径。 | ||
| 6 | + | ||
| 7 | +## 接口 | ||
| 8 | + | ||
| 9 | +| 方法与路径 | Envelope `type` | 功能 | | ||
| 10 | +| --- | --- | --- | | ||
| 11 | +| `POST /api/db/write` | `db/write` | 按 ID 创建或更新数据库记录 | | ||
| 12 | +| `POST /api/db/read` | `db/read` | 按 ID 读取数据库记录 | | ||
| 13 | +| `POST /api/redis/write` | `redis/write` | 写入带 TTL 的 Redis 字符串 | | ||
| 14 | +| `POST /api/redis/read` | `redis/read` | 读取 Redis 字符串 | | ||
| 15 | +| `POST /api/envelope/inspect` | `envelope/inspect` | 返回 Envelope、Metadata 和 RequestContext 解析结果 | | ||
| 16 | +| `POST /api/k8s/pod/read` | `k8s/pod/read` | 查询 Demo 管理的 Pod | | ||
| 17 | +| `POST /api/k8s/pod/create` | `k8s/pod/create` | 使用固定模板创建 Pod | | ||
| 18 | +| `POST /api/k8s/pod/delete` | `k8s/pod/delete` | 删除 Demo 管理的 Pod | | ||
| 19 | + | ||
| 20 | +Swagger UI 位于 `/docs`。OpenAPI 文档位于 `/openapi.json`。 | ||
| 21 | + | ||
| 22 | +## 本地启动 | ||
| 23 | + | ||
| 24 | +在 `service` 目录执行: | ||
| 25 | + | ||
| 26 | +```powershell | ||
| 27 | +uv sync | ||
| 28 | +Copy-Item ` | ||
| 29 | + -LiteralPath ".\examples\simple_capabilities.local.env.example" ` | ||
| 30 | + -Destination ".\examples\.env.development.local" | ||
| 31 | +uv run uvicorn examples.simple_capabilities_app:asgi ` | ||
| 32 | + --host 127.0.0.1 ` | ||
| 33 | + --port 8090 ` | ||
| 34 | + --workers 1 ` | ||
| 35 | + --env-file examples/.env.development.local | ||
| 36 | +``` | ||
| 37 | + | ||
| 38 | +浏览器访问 <http://127.0.0.1:8090/docs>。本地 SQLite 数据保存在 | ||
| 39 | +`service/simple_capabilities_demo.db`。FakeRedis 和 Fake Kubernetes 数据保存在当前 Python 进程中, | ||
| 40 | +服务重启后清空。 | ||
| 41 | + | ||
| 42 | +## Swagger 调试 | ||
| 43 | + | ||
| 44 | +所有接口接收完整 Envelope。`metadata.request_id` 为必填字段。请求路径与 `type` 保持一致。 | ||
| 45 | + | ||
| 46 | +### 写数据库 | ||
| 47 | + | ||
| 48 | +```json | ||
| 49 | +{ | ||
| 50 | + "type": "db/write", | ||
| 51 | + "metadata": {"request_id": "db-write-1"}, | ||
| 52 | + "rawdata": { | ||
| 53 | + "id": "record-1", | ||
| 54 | + "value": "hello database" | ||
| 55 | + }, | ||
| 56 | + "version": "1" | ||
| 57 | +} | ||
| 58 | +``` | ||
| 59 | + | ||
| 60 | +首次写入返回 `operation=created`。使用新 `request_id` 再次写入同一 ID 返回 | ||
| 61 | +`operation=updated`。 | ||
| 62 | + | ||
| 63 | +### 读数据库 | ||
| 64 | + | ||
| 65 | +```json | ||
| 66 | +{ | ||
| 67 | + "type": "db/read", | ||
| 68 | + "metadata": {"request_id": "db-read-1"}, | ||
| 69 | + "rawdata": {"id": "record-1"}, | ||
| 70 | + "version": "1" | ||
| 71 | +} | ||
| 72 | +``` | ||
| 73 | + | ||
| 74 | +缺失记录返回 HTTP `404` 和 `error_code=not_found`。 | ||
| 75 | + | ||
| 76 | +### 写 Redis | ||
| 77 | + | ||
| 78 | +```json | ||
| 79 | +{ | ||
| 80 | + "type": "redis/write", | ||
| 81 | + "metadata": {"request_id": "redis-write-1"}, | ||
| 82 | + "rawdata": { | ||
| 83 | + "key": "message-1", | ||
| 84 | + "value": "hello redis", | ||
| 85 | + "ttl_seconds": 300 | ||
| 86 | + }, | ||
| 87 | + "version": "1" | ||
| 88 | +} | ||
| 89 | +``` | ||
| 90 | + | ||
| 91 | +省略 `ttl_seconds` 时使用 `DEMO_REDIS_DEFAULT_TTL_SECONDS`。实际 Redis key 为 | ||
| 92 | +`{OPENJIUWEN_SERVICE_REDIS_KEY_PREFIX}:kv:{key}`。 | ||
| 93 | + | ||
| 94 | +### 读 Redis | ||
| 95 | + | ||
| 96 | +```json | ||
| 97 | +{ | ||
| 98 | + "type": "redis/read", | ||
| 99 | + "metadata": {"request_id": "redis-read-1"}, | ||
| 100 | + "rawdata": {"key": "message-1"}, | ||
| 101 | + "version": "1" | ||
| 102 | +} | ||
| 103 | +``` | ||
| 104 | + | ||
| 105 | +缺失 key 返回 HTTP `200`、`found=false` 和 `value=null`。 | ||
| 106 | + | ||
| 107 | +### 解析 Envelope | ||
| 108 | + | ||
| 109 | +```json | ||
| 110 | +{ | ||
| 111 | + "type": "envelope/inspect", | ||
| 112 | + "metadata": { | ||
| 113 | + "request_id": "inspect-1", | ||
| 114 | + "user_id": "demo-user", | ||
| 115 | + "chat_id": "demo-chat", | ||
| 116 | + "session_id": "demo-session", | ||
| 117 | + "bot_id": "demo-bot", | ||
| 118 | + "channel": "swagger", | ||
| 119 | + "timestamp": 1786500000, | ||
| 120 | + "trace_id": "trace-1", | ||
| 121 | + "instance_id": "demo-instance", | ||
| 122 | + "extra": {"tenant": "demo"} | ||
| 123 | + }, | ||
| 124 | + "rawdata": { | ||
| 125 | + "message": "inspect this envelope", | ||
| 126 | + "attributes": {"source": "swagger"} | ||
| 127 | + }, | ||
| 128 | + "version": "1" | ||
| 129 | +} | ||
| 130 | +``` | ||
| 131 | + | ||
| 132 | +响应中的 `envelope` 来自框架解析后的 Envelope,`context` 来自当前 RequestContext,`service` | ||
| 133 | +包含当前 Demo 环境标识、服务标签和 Redis 模式。 | ||
| 134 | + | ||
| 135 | +### 创建 Pod | ||
| 136 | + | ||
| 137 | +```json | ||
| 138 | +{ | ||
| 139 | + "type": "k8s/pod/create", | ||
| 140 | + "metadata": {"request_id": "pod-create-1"}, | ||
| 141 | + "rawdata": {"name": "capability-pod-1"}, | ||
| 142 | + "version": "1" | ||
| 143 | +} | ||
| 144 | +``` | ||
| 145 | + | ||
| 146 | +服务使用 `DEMO_KUBERNETES_POD_IMAGE` 和固定安全模板创建 Pod。请求体只接收符合 Kubernetes | ||
| 147 | +DNS-1123 subdomain 规则的名称。同名 Pod 返回 HTTP `409` 和 `error_code=conflict`。 | ||
| 148 | + | ||
| 149 | +### 查询 Pod | ||
| 150 | + | ||
| 151 | +```json | ||
| 152 | +{ | ||
| 153 | + "type": "k8s/pod/read", | ||
| 154 | + "metadata": {"request_id": "pod-read-1"}, | ||
| 155 | + "rawdata": {"name": "capability-pod-1"}, | ||
| 156 | + "version": "1" | ||
| 157 | +} | ||
| 158 | +``` | ||
| 159 | + | ||
| 160 | +响应包含 `name`、`namespace`、`phase`、`ready` 和 `image`。Pod 缺失或管理标签不匹配时返回 HTTP | ||
| 161 | +`404` 和 `error_code=not_found`。 | ||
| 162 | + | ||
| 163 | +### 删除 Pod | ||
| 164 | + | ||
| 165 | +```json | ||
| 166 | +{ | ||
| 167 | + "type": "k8s/pod/delete", | ||
| 168 | + "metadata": {"request_id": "pod-delete-1"}, | ||
| 169 | + "rawdata": {"name": "capability-pod-1"}, | ||
| 170 | + "version": "1" | ||
| 171 | +} | ||
| 172 | +``` | ||
| 173 | + | ||
| 174 | +删除结果的 `state` 为 `delete_requested`、`deletion_in_progress` 或 `already_absent`。Real 模式使用 | ||
| 175 | +读取结果中的 Pod UID 提交删除前置条件。 | ||
| 176 | + | ||
| 177 | +## 配置 | ||
| 178 | + | ||
| 179 | +Demo 自有变量: | ||
| 180 | + | ||
| 181 | +| 变量 | 说明 | | ||
| 182 | +| --- | --- | | ||
| 183 | +| `DEMO_ENVIRONMENT` | 运行环境标识 | | ||
| 184 | +| `DEMO_SERVICE_LABEL` | 服务实例标签 | | ||
| 185 | +| `DEMO_REDIS_MODE` | `fake` 或 `real`,必须显式配置 | | ||
| 186 | +| `DEMO_REDIS_DEFAULT_TTL_SECONDS` | Redis 默认 TTL,取正整数 | | ||
| 187 | +| `DEMO_KUBERNETES_MODE` | `fake` 或 `real`,必须显式配置 | | ||
| 188 | +| `DEMO_KUBERNETES_NAMESPACE` | Pod 操作使用的固定 namespace,符合 DNS-1123 label 规则 | | ||
| 189 | +| `DEMO_KUBERNETES_POD_IMAGE` | Pod 固定模板使用的镜像 | | ||
| 190 | + | ||
| 191 | +`fake` 模式要求 `OPENJIUWEN_SERVICE_REDIS_URL=disabled` 或 `none`。`real` 模式要求有效的 | ||
| 192 | +`redis://`、`rediss://` 或 `unix://` URL。数据库、锁、缓存和 Redis 命名空间沿用 | ||
| 193 | +`OPENJIUWEN_SERVICE_*` 配置。Real Kubernetes 模式优先读取 in-cluster 配置,随后读取 `KUBECONFIG` | ||
| 194 | +或默认 kubeconfig。 | ||
| 195 | + | ||
| 196 | +## 服务器部署 | ||
| 197 | + | ||
| 198 | +从模板生成实际配置文件: | ||
| 199 | + | ||
| 200 | +```powershell | ||
| 201 | +Copy-Item ` | ||
| 202 | + -LiteralPath ".\examples\simple_capabilities.server.env.example" ` | ||
| 203 | + -Destination ".\examples\.env.production.local" | ||
| 204 | +``` | ||
| 205 | + | ||
| 206 | +编辑实际配置文件中的 MySQL 地址、数据库名、账号、密码、Redis URL、Kubernetes namespace 和镜像。 | ||
| 207 | +MySQL 账号需要具备连接、建表和 CRUD 权限。安装 Real Kubernetes 后端并启动服务: | ||
| 208 | + | ||
| 209 | +```powershell | ||
| 210 | +uv sync --extra kubernetes | ||
| 211 | +uv run --extra kubernetes uvicorn examples.simple_capabilities_app:asgi ` | ||
| 212 | + --host 0.0.0.0 ` | ||
| 213 | + --port 8090 ` | ||
| 214 | + --workers 1 ` | ||
| 215 | + --env-file examples/.env.production.local | ||
| 216 | +``` | ||
| 217 | + | ||
| 218 | +Nginx 或 API Gateway 提供 HTTPS、访问来源限制和请求日志。应用启动时连接数据库与 Redis、创建 | ||
| 219 | +`simple_capability_records` 表并执行资源探活。Kubernetes 探活通过 namespace Pod list 验证连接和 | ||
| 220 | +权限。资源连接或配置校验失败会终止启动。 | ||
| 221 | + | ||
| 222 | +Kubernetes ServiceAccount 在目标 namespace 中需要以下 Pod 权限: | ||
| 223 | + | ||
| 224 | +```yaml | ||
| 225 | +apiVersion: rbac.authorization.k8s.io/v1 | ||
| 226 | +kind: Role | ||
| 227 | +metadata: | ||
| 228 | + name: simple-capabilities-demo-pods | ||
| 229 | + namespace: simple-capabilities-demo | ||
| 230 | +rules: | ||
| 231 | + - apiGroups: [""] | ||
| 232 | + resources: ["pods"] | ||
| 233 | + verbs: ["get", "list", "create", "delete"] | ||
| 234 | +``` | ||
| 235 | + | ||
| 236 | +RoleBinding 将该 Role 绑定到运行 Demo 服务的 ServiceAccount。查询和删除操作要求 Pod 包含 | ||
| 237 | +`app.kubernetes.io/managed-by=openjiuwen-service-demo` 标签。 | ||
| 238 | + | ||
| 239 | +## 测试 | ||
| 240 | + | ||
| 241 | +在 `service` 目录执行专项测试: | ||
| 242 | + | ||
| 243 | +```powershell | ||
| 244 | +uv run pytest ` | ||
| 245 | + tests/unit_tests/test_kubernetes_operations.py ` | ||
| 246 | + tests/unit_tests/test_system_context.py ` | ||
| 247 | + tests/unit_tests/test_config_bootstrap.py ` | ||
| 248 | + tests/unit_tests/test_simple_capabilities_example.py ` | ||
| 249 | + -q -p no:cacheprovider | ||
| 250 | +``` | ||
| 251 | + | ||
| 252 | +执行包含 Real 适配器的专项测试: | ||
| 253 | + | ||
| 254 | +```powershell | ||
| 255 | +uv run --extra kubernetes pytest ` | ||
| 256 | + tests/unit_tests/test_kubernetes_operations.py ` | ||
| 257 | + -q -p no:cacheprovider | ||
| 258 | +``` | ||
| 259 | + | ||
| 260 | +执行 Service 全量测试: | ||
| 261 | + | ||
| 262 | +```powershell | ||
| 263 | +uv run pytest tests -q -p no:cacheprovider | ||
| 264 | +``` | ||
| @@ -0,0 +1,18 @@ | |||
| 1 | +DEMO_ENVIRONMENT=local | ||
| 2 | +DEMO_SERVICE_LABEL=simple-capabilities-demo | ||
| 3 | +DEMO_REDIS_MODE=fake | ||
| 4 | +DEMO_REDIS_DEFAULT_TTL_SECONDS=300 | ||
| 5 | +DEMO_KUBERNETES_MODE=fake | ||
| 6 | +DEMO_KUBERNETES_NAMESPACE=simple-capabilities-demo | ||
| 7 | +DEMO_KUBERNETES_POD_IMAGE=demo-pod:local | ||
| 8 | + | ||
| 9 | +OPENJIUWEN_SERVICE_DB_TYPE=sqlite | ||
| 10 | +OPENJIUWEN_SERVICE_DB_NAME=./simple_capabilities_demo.db | ||
| 11 | + | ||
| 12 | +OPENJIUWEN_SERVICE_REDIS_URL=disabled | ||
| 13 | +OPENJIUWEN_SERVICE_REDIS_KEY_PREFIX=simple-capabilities-demo | ||
| 14 | + | ||
| 15 | +OPENJIUWEN_SERVICE_LOCK_BACKEND=memory | ||
| 16 | +OPENJIUWEN_SERVICE_CACHE_BACKEND=memory | ||
| 17 | +OPENJIUWEN_SERVICE_DEPLOY_REPLICAS=1 | ||
| 18 | +OPENJIUWEN_SERVICE_REQUEST_TIMEOUT_SECONDS=30 | ||
| @@ -0,0 +1,22 @@ | |||
| 1 | +DEMO_ENVIRONMENT=server | ||
| 2 | +DEMO_SERVICE_LABEL=simple-capabilities-demo | ||
| 3 | +DEMO_REDIS_MODE=real | ||
| 4 | +DEMO_REDIS_DEFAULT_TTL_SECONDS=300 | ||
| 5 | +DEMO_KUBERNETES_MODE=real | ||
| 6 | +DEMO_KUBERNETES_NAMESPACE=simple-capabilities-demo | ||
| 7 | +DEMO_KUBERNETES_POD_IMAGE=registry.example.com/simple-capabilities/pod-demo:1.0.0 | ||
| 8 | + | ||
| 9 | +OPENJIUWEN_SERVICE_DB_TYPE=mysql | ||
| 10 | +OPENJIUWEN_SERVICE_DB_HOST=127.0.0.1 | ||
| 11 | +OPENJIUWEN_SERVICE_DB_PORT=3306 | ||
| 12 | +OPENJIUWEN_SERVICE_DB_NAME=simple_capabilities_demo | ||
| 13 | +OPENJIUWEN_SERVICE_DB_USER=simple_capabilities_demo | ||
| 14 | +OPENJIUWEN_SERVICE_DB_PASSWORD=change-me | ||
| 15 | + | ||
| 16 | +OPENJIUWEN_SERVICE_REDIS_URL=redis://127.0.0.1:6379/0 | ||
| 17 | +OPENJIUWEN_SERVICE_REDIS_KEY_PREFIX=simple-capabilities-demo | ||
| 18 | + | ||
| 19 | +OPENJIUWEN_SERVICE_LOCK_BACKEND=redis | ||
| 20 | +OPENJIUWEN_SERVICE_CACHE_BACKEND=memory | ||
| 21 | +OPENJIUWEN_SERVICE_DEPLOY_REPLICAS=1 | ||
| 22 | +OPENJIUWEN_SERVICE_REQUEST_TIMEOUT_SECONDS=30 | ||
| @@ -0,0 +1,567 @@ | |||
| 1 | +# coding: utf-8 | ||
| 2 | +# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved | ||
| 3 | + | ||
| 4 | +"""Deployable demo for database, Redis, Envelope, and Kubernetes capabilities.""" | ||
| 5 | + | ||
| 6 | +from __future__ import annotations | ||
| 7 | + | ||
| 8 | +import os | ||
| 9 | +import re | ||
| 10 | +from collections.abc import Mapping | ||
| 11 | +from dataclasses import asdict, dataclass | ||
| 12 | +from datetime import datetime, timezone | ||
| 13 | +from typing import Any, Literal | ||
| 14 | +from urllib.parse import urlparse | ||
| 15 | + | ||
| 16 | +from pydantic import BaseModel, ConfigDict, Field, field_serializer | ||
| 17 | + | ||
| 18 | +from openjiuwen_runtime.foundation.db.table_def import ( | ||
| 19 | + ColumnDefinition, | ||
| 20 | + TableDefinition, | ||
| 21 | +) | ||
| 22 | +from openjiuwen_runtime.service import ( | ||
| 23 | + App, | ||
| 24 | + DatabaseUnavailable, | ||
| 25 | + Envelope, | ||
| 26 | + FakeKubernetesOperations, | ||
| 27 | + KubernetesAsyncioOperations, | ||
| 28 | + NotFoundError, | ||
| 29 | + PodCreateSpec, | ||
| 30 | + RedisUnavailable, | ||
| 31 | + ServiceConfig, | ||
| 32 | + SystemContext, | ||
| 33 | + TypedAppContext, | ||
| 34 | + build_system_context, | ||
| 35 | +) | ||
| 36 | + | ||
| 37 | + | ||
| 38 | +DEMO_TABLE_NAME = "simple_capability_records" | ||
| 39 | +DEMO_TABLE = TableDefinition( | ||
| 40 | + table_name=DEMO_TABLE_NAME, | ||
| 41 | + columns=[ | ||
| 42 | + ColumnDefinition("id", "string", primary_key=True, nullable=False, length=64), | ||
| 43 | + ColumnDefinition("value", "string", nullable=False, length=4096), | ||
| 44 | + ColumnDefinition("updated_at", "datetime", nullable=False), | ||
| 45 | + ], | ||
| 46 | +) | ||
| 47 | + | ||
| 48 | +_DISABLED_REDIS_URLS = {"disabled", "none"} | ||
| 49 | +_DNS_1123_LABEL_PATTERN = r"^[a-z0-9](?:[-a-z0-9]{0,61}[a-z0-9])?$" | ||
| 50 | +_DNS_1123_SUBDOMAIN_PATTERN = ( | ||
| 51 | + r"^[a-z0-9](?:[-a-z0-9]{0,61}[a-z0-9])?" | ||
| 52 | + r"(?:\.[a-z0-9](?:[-a-z0-9]{0,61}[a-z0-9])?)*$" | ||
| 53 | +) | ||
| 54 | + | ||
| 55 | + | ||
| 56 | + | ||
| 57 | +class DemoConfig: | ||
| 58 | + """Configuration owned by the simple capabilities demo.""" | ||
| 59 | + | ||
| 60 | + environment: str | ||
| 61 | + service_label: str | ||
| 62 | + redis_mode: Literal["fake", "real"] | ||
| 63 | + redis_default_ttl_seconds: int | ||
| 64 | + kubernetes_mode: Literal["fake", "real"] | ||
| 65 | + kubernetes_namespace: str | ||
| 66 | + kubernetes_pod_image: str | ||
| 67 | + | ||
| 68 | + def __post_init__(self) -> None: | ||
| 69 | + environment = str(self.environment).strip() | ||
| 70 | + service_label = str(self.service_label).strip() | ||
| 71 | + redis_mode = str(self.redis_mode).strip().lower() | ||
| 72 | + kubernetes_mode = str(self.kubernetes_mode).strip().lower() | ||
| 73 | + kubernetes_namespace = str(self.kubernetes_namespace).strip() | ||
| 74 | + kubernetes_pod_image = str(self.kubernetes_pod_image).strip() | ||
| 75 | + if not environment: | ||
| 76 | + raise ValueError("DEMO_ENVIRONMENT must not be empty") | ||
| 77 | + if not service_label: | ||
| 78 | + raise ValueError("DEMO_SERVICE_LABEL must not be empty") | ||
| 79 | + if redis_mode not in {"fake", "real"}: | ||
| 80 | + raise ValueError("DEMO_REDIS_MODE must be one of fake, real") | ||
| 81 | + if kubernetes_mode not in {"fake", "real"}: | ||
| 82 | + raise ValueError("DEMO_KUBERNETES_MODE must be one of fake, real") | ||
| 83 | + if not re.fullmatch(_DNS_1123_LABEL_PATTERN, kubernetes_namespace): | ||
| 84 | + raise ValueError( | ||
| 85 | + "DEMO_KUBERNETES_NAMESPACE must be a valid DNS-1123 label" | ||
| 86 | + ) | ||
| 87 | + if not kubernetes_pod_image: | ||
| 88 | + raise ValueError("DEMO_KUBERNETES_POD_IMAGE must not be empty") | ||
| 89 | + ttl = self.redis_default_ttl_seconds | ||
| 90 | + if isinstance(ttl, bool) or not isinstance(ttl, int) or ttl <= 0: | ||
| 91 | + raise ValueError( | ||
| 92 | + "DEMO_REDIS_DEFAULT_TTL_SECONDS must be a positive integer" | ||
| 93 | + ) | ||
| 94 | + object.__setattr__(self, "environment", environment) | ||
| 95 | + object.__setattr__(self, "service_label", service_label) | ||
| 96 | + object.__setattr__(self, "redis_mode", redis_mode) | ||
| 97 | + object.__setattr__(self, "kubernetes_mode", kubernetes_mode) | ||
| 98 | + object.__setattr__(self, "kubernetes_namespace", kubernetes_namespace) | ||
| 99 | + object.__setattr__(self, "kubernetes_pod_image", kubernetes_pod_image) | ||
| 100 | + | ||
| 101 | + | ||
| 102 | + def from_env(cls) -> "DemoConfig": | ||
| 103 | + redis_mode = os.getenv("DEMO_REDIS_MODE") | ||
| 104 | + if redis_mode is None or not redis_mode.strip(): | ||
| 105 | + raise ValueError("DEMO_REDIS_MODE is required") | ||
| 106 | + kubernetes_mode = os.getenv("DEMO_KUBERNETES_MODE") | ||
| 107 | + if kubernetes_mode is None or not kubernetes_mode.strip(): | ||
| 108 | + raise ValueError("DEMO_KUBERNETES_MODE is required") | ||
| 109 | + try: | ||
| 110 | + ttl = int(os.getenv("DEMO_REDIS_DEFAULT_TTL_SECONDS", "300")) | ||
| 111 | + except ValueError as exc: | ||
| 112 | + raise ValueError( | ||
| 113 | + "DEMO_REDIS_DEFAULT_TTL_SECONDS must be a positive integer" | ||
| 114 | + ) from exc | ||
| 115 | + config = cls( | ||
| 116 | + environment=os.getenv("DEMO_ENVIRONMENT", ""), | ||
| 117 | + service_label=os.getenv("DEMO_SERVICE_LABEL", ""), | ||
| 118 | + redis_mode=redis_mode, | ||
| 119 | + redis_default_ttl_seconds=ttl, | ||
| 120 | + kubernetes_mode=kubernetes_mode, | ||
| 121 | + kubernetes_namespace=os.getenv("DEMO_KUBERNETES_NAMESPACE", ""), | ||
| 122 | + kubernetes_pod_image=os.getenv("DEMO_KUBERNETES_POD_IMAGE", ""), | ||
| 123 | + ) | ||
| 124 | + config.validate_redis_url( | ||
| 125 | + os.getenv( | ||
| 126 | + "OPENJIUWEN_SERVICE_REDIS_URL", | ||
| 127 | + "redis://localhost:6379/0", | ||
| 128 | + ) | ||
| 129 | + ) | ||
| 130 | + return config | ||
| 131 | + | ||
| 132 | + def validate_redis_url(self, redis_url: str | None) -> None: | ||
| 133 | + value = (redis_url or "").strip() | ||
| 134 | + if self.redis_mode == "fake": | ||
| 135 | + if value.lower() not in _DISABLED_REDIS_URLS: | ||
| 136 | + raise ValueError( | ||
| 137 | + "fake Redis mode requires OPENJIUWEN_SERVICE_REDIS_URL=" | ||
| 138 | + "disabled or none" | ||
| 139 | + ) | ||
| 140 | + return | ||
| 141 | + | ||
| 142 | + if value.lower() in _DISABLED_REDIS_URLS or not value: | ||
| 143 | + raise ValueError("real Redis mode requires OPENJIUWEN_SERVICE_REDIS_URL") | ||
| 144 | + parsed = urlparse(value) | ||
| 145 | + network_url = parsed.scheme in {"redis", "rediss"} and parsed.hostname | ||
| 146 | + unix_url = parsed.scheme == "unix" and parsed.path | ||
| 147 | + if not (network_url or unix_url): | ||
| 148 | + raise ValueError("OPENJIUWEN_SERVICE_REDIS_URL must be a valid Redis URL") | ||
| 149 | + | ||
| 150 | + | ||
| 151 | +class _StrictInput(BaseModel): | ||
| 152 | + model_config = ConfigDict(extra="forbid", str_strip_whitespace=True) | ||
| 153 | + | ||
| 154 | + | ||
| 155 | +class DbWriteInput(_StrictInput): | ||
| 156 | + id: str = Field(min_length=1, max_length=64) | ||
| 157 | + value: str = Field(min_length=1, max_length=4096) | ||
| 158 | + | ||
| 159 | + | ||
| 160 | +class DbReadInput(_StrictInput): | ||
| 161 | + id: str = Field(min_length=1, max_length=64) | ||
| 162 | + | ||
| 163 | + | ||
| 164 | +class RedisWriteInput(_StrictInput): | ||
| 165 | + key: str = Field(min_length=1, max_length=256) | ||
| 166 | + value: str = Field(min_length=1, max_length=4096) | ||
| 167 | + ttl_seconds: int | None = Field(default=None, gt=0) | ||
| 168 | + | ||
| 169 | + | ||
| 170 | +class RedisReadInput(_StrictInput): | ||
| 171 | + key: str = Field(min_length=1, max_length=256) | ||
| 172 | + | ||
| 173 | + | ||
| 174 | +class EnvelopeInspectInput(_StrictInput): | ||
| 175 | + message: str = Field(min_length=1, max_length=4096) | ||
| 176 | + attributes: dict[str, Any] = Field(default_factory=dict) | ||
| 177 | + | ||
| 178 | + | ||
| 179 | +class PodNameInput(_StrictInput): | ||
| 180 | + name: str = Field( | ||
| 181 | + min_length=1, | ||
| 182 | + max_length=253, | ||
| 183 | + pattern=_DNS_1123_SUBDOMAIN_PATTERN, | ||
| 184 | + ) | ||
| 185 | + | ||
| 186 | + | ||
| 187 | +class RecordView(BaseModel): | ||
| 188 | + id: str | ||
| 189 | + value: str | ||
| 190 | + updated_at: datetime | ||
| 191 | + | ||
| 192 | + | ||
| 193 | + def serialize_updated_at(self, value: datetime) -> str: | ||
| 194 | + if value.tzinfo is None: | ||
| 195 | + value = value.replace(tzinfo=timezone.utc) | ||
| 196 | + return value.astimezone(timezone.utc).isoformat().replace("+00:00", "Z") | ||
| 197 | + | ||
| 198 | + | ||
| 199 | +class DbWriteOutput(BaseModel): | ||
| 200 | + operation: Literal["created", "updated"] | ||
| 201 | + record: RecordView | ||
| 202 | + | ||
| 203 | + | ||
| 204 | +class DbReadOutput(BaseModel): | ||
| 205 | + record: RecordView | ||
| 206 | + | ||
| 207 | + | ||
| 208 | +class RedisWriteOutput(BaseModel): | ||
| 209 | + key: str | ||
| 210 | + value: str | ||
| 211 | + ttl_seconds: int | ||
| 212 | + backend_mode: Literal["fake", "real"] | ||
| 213 | + | ||
| 214 | + | ||
| 215 | +class RedisReadOutput(BaseModel): | ||
| 216 | + key: str | ||
| 217 | + found: bool | ||
| 218 | + value: str | None | ||
| 219 | + backend_mode: Literal["fake", "real"] | ||
| 220 | + | ||
| 221 | + | ||
| 222 | +class MetadataView(BaseModel): | ||
| 223 | + request_id: str | ||
| 224 | + user_id: str | None | ||
| 225 | + chat_id: str | None | ||
| 226 | + session_id: str | None | ||
| 227 | + bot_id: str | None | ||
| 228 | + channel: str | None | ||
| 229 | + timestamp: float | None | ||
| 230 | + trace_id: str | None | ||
| 231 | + instance_id: str | None | ||
| 232 | + extra: dict[str, Any] | ||
| 233 | + | ||
| 234 | + | ||
| 235 | +class InspectedEnvelopeView(BaseModel): | ||
| 236 | + type: str | ||
| 237 | + version: str | ||
| 238 | + metadata: MetadataView | ||
| 239 | + rawdata: EnvelopeInspectInput | ||
| 240 | + | ||
| 241 | + | ||
| 242 | +class RequestContextView(BaseModel): | ||
| 243 | + msg_type: str | ||
| 244 | + request_id: str | ||
| 245 | + user_id: str | None | ||
| 246 | + chat_id: str | None | ||
| 247 | + session_id: str | None | ||
| 248 | + bot_id: str | None | ||
| 249 | + channel: str | None | ||
| 250 | + trace_id: str | None | ||
| 251 | + instance_id: str | None | ||
| 252 | + replica_id: str | ||
| 253 | + | ||
| 254 | + | ||
| 255 | +class ServiceView(BaseModel): | ||
| 256 | + environment: str | ||
| 257 | + service_label: str | ||
| 258 | + redis_mode: Literal["fake", "real"] | ||
| 259 | + | ||
| 260 | + | ||
| 261 | +class EnvelopeInspectOutput(BaseModel): | ||
| 262 | + envelope: InspectedEnvelopeView | ||
| 263 | + context: RequestContextView | ||
| 264 | + service: ServiceView | ||
| 265 | + | ||
| 266 | + | ||
| 267 | +class PodView(BaseModel): | ||
| 268 | + name: str | ||
| 269 | + namespace: str | ||
| 270 | + phase: str | ||
| 271 | + ready: bool | ||
| 272 | + image: str | None | ||
| 273 | + | ||
| 274 | + | ||
| 275 | +class PodReadOutput(BaseModel): | ||
| 276 | + pod: PodView | ||
| 277 | + | ||
| 278 | + | ||
| 279 | +class PodCreateOutput(BaseModel): | ||
| 280 | + operation: Literal["created"] | ||
| 281 | + pod: PodView | ||
| 282 | + | ||
| 283 | + | ||
| 284 | +class PodDeleteOutput(BaseModel): | ||
| 285 | + name: str | ||
| 286 | + namespace: str | ||
| 287 | + state: Literal[ | ||
| 288 | + "delete_requested", | ||
| 289 | + "deletion_in_progress", | ||
| 290 | + "already_absent", | ||
| 291 | + ] | ||
| 292 | + | ||
| 293 | + | ||
| 294 | +def build_demo_system_context( | ||
| 295 | + service_config: ServiceConfig, | ||
| 296 | + demo_config: DemoConfig, | ||
| 297 | +) -> SystemContext: | ||
| 298 | + """Build resources for the selected local or server deployment mode.""" | ||
| 299 | + demo_config.validate_redis_url(service_config.redis_url) | ||
| 300 | + redis_client = None | ||
| 301 | + if demo_config.redis_mode == "fake": | ||
| 302 | + import fakeredis.aioredis | ||
| 303 | + | ||
| 304 | + redis_client = fakeredis.aioredis.FakeRedis() | ||
| 305 | + | ||
| 306 | + if demo_config.kubernetes_mode == "fake": | ||
| 307 | + kubernetes = FakeKubernetesOperations(demo_config.kubernetes_namespace) | ||
| 308 | + else: | ||
| 309 | + kubernetes = KubernetesAsyncioOperations( | ||
| 310 | + demo_config.kubernetes_namespace, | ||
| 311 | + kubeconfig=os.getenv("KUBECONFIG"), | ||
| 312 | + ) | ||
| 313 | + | ||
| 314 | + system = build_system_context( | ||
| 315 | + service_config, | ||
| 316 | + redis=redis_client, | ||
| 317 | + kubernetes=kubernetes, | ||
| 318 | + table_definitions=(DEMO_TABLE,), | ||
| 319 | + ) | ||
| 320 | + if redis_client is not None: | ||
| 321 | + system.set_redis(redis_client, owned=True) | ||
| 322 | + system.set_kubernetes(kubernetes, owned=True) | ||
| 323 | + return system | ||
| 324 | + | ||
| 325 | + | ||
| 326 | +def _record_view(record: Any) -> RecordView: | ||
| 327 | + if isinstance(record, Mapping): | ||
| 328 | + values = record | ||
| 329 | + else: | ||
| 330 | + to_dict = getattr(record, "to_dict", None) | ||
| 331 | + if callable(to_dict): | ||
| 332 | + values = to_dict() | ||
| 333 | + else: | ||
| 334 | + values = { | ||
| 335 | + "id": record.id, | ||
| 336 | + "value": record.value, | ||
| 337 | + "updated_at": record.updated_at, | ||
| 338 | + } | ||
| 339 | + return RecordView.model_validate(values) | ||
| 340 | + | ||
| 341 | + | ||
| 342 | +def register_handlers(application: App, config: DemoConfig) -> None: | ||
| 343 | + | ||
| 344 | + "db/write", | ||
| 345 | + request_model=DbWriteInput, | ||
| 346 | + response_model=DbWriteOutput, | ||
| 347 | + summary="Write a database record", | ||
| 348 | + description="Create or update one record by ID.", | ||
| 349 | + tags=["Database"], | ||
| 350 | + ) | ||
| 351 | + async def write_database( | ||
| 352 | + ctx: TypedAppContext[DbWriteInput], | ||
| 353 | + env: Envelope[DbWriteInput], | ||
| 354 | + ) -> dict[str, Any]: | ||
| 355 | + request = ctx.request | ||
| 356 | + updated_at = datetime.now(timezone.utc) | ||
| 357 | + try: | ||
| 358 | + existing = await ctx.db_get(DEMO_TABLE_NAME, {"id": request.id}) | ||
| 359 | + if existing is None: | ||
| 360 | + operation = "created" | ||
| 361 | + await ctx.db_create( | ||
| 362 | + DEMO_TABLE_NAME, | ||
| 363 | + { | ||
| 364 | + "id": request.id, | ||
| 365 | + "value": request.value, | ||
| 366 | + "updated_at": updated_at, | ||
| 367 | + }, | ||
| 368 | + ) | ||
| 369 | + else: | ||
| 370 | + operation = "updated" | ||
| 371 | + await ctx.db_update( | ||
| 372 | + DEMO_TABLE_NAME, | ||
| 373 | + {"id": request.id}, | ||
| 374 | + {"value": request.value, "updated_at": updated_at}, | ||
| 375 | + ) | ||
| 376 | + record = await ctx.db_get(DEMO_TABLE_NAME, {"id": request.id}) | ||
| 377 | + except DatabaseUnavailable: | ||
| 378 | + raise | ||
| 379 | + except Exception as exc: | ||
| 380 | + raise DatabaseUnavailable("database operation failed") from exc | ||
| 381 | + if record is None: | ||
| 382 | + raise DatabaseUnavailable("database write did not return a record") | ||
| 383 | + return { | ||
| 384 | + "operation": operation, | ||
| 385 | + "record": _record_view(record).model_dump(mode="python"), | ||
| 386 | + } | ||
| 387 | + | ||
| 388 | + | ||
| 389 | + "db/read", | ||
| 390 | + request_model=DbReadInput, | ||
| 391 | + response_model=DbReadOutput, | ||
| 392 | + summary="Read a database record", | ||
| 393 | + description="Read one record by ID.", | ||
| 394 | + tags=["Database"], | ||
| 395 | + ) | ||
| 396 | + async def read_database( | ||
| 397 | + ctx: TypedAppContext[DbReadInput], | ||
| 398 | + env: Envelope[DbReadInput], | ||
| 399 | + ) -> dict[str, Any]: | ||
| 400 | + try: | ||
| 401 | + record = await ctx.db_get(DEMO_TABLE_NAME, {"id": ctx.request.id}) | ||
| 402 | + except DatabaseUnavailable: | ||
| 403 | + raise | ||
| 404 | + except Exception as exc: | ||
| 405 | + raise DatabaseUnavailable("database operation failed") from exc | ||
| 406 | + if record is None: | ||
| 407 | + raise NotFoundError(f"record {ctx.request.id!r} not found") | ||
| 408 | + return {"record": _record_view(record).model_dump(mode="python")} | ||
| 409 | + | ||
| 410 | + | ||
| 411 | + "redis/write", | ||
| 412 | + request_model=RedisWriteInput, | ||
| 413 | + response_model=RedisWriteOutput, | ||
| 414 | + summary="Write a Redis value", | ||
| 415 | + description="Write one namespaced Redis key with a TTL.", | ||
| 416 | + tags=["Redis"], | ||
| 417 | + ) | ||
| 418 | + async def write_redis( | ||
| 419 | + ctx: TypedAppContext[RedisWriteInput], | ||
| 420 | + env: Envelope[RedisWriteInput], | ||
| 421 | + ) -> dict[str, Any]: | ||
| 422 | + request = ctx.request | ||
| 423 | + ttl = request.ttl_seconds or config.redis_default_ttl_seconds | ||
| 424 | + try: | ||
| 425 | + await ctx.kv.set(request.key, request.value, ttl=ttl) | ||
| 426 | + except RedisUnavailable: | ||
| 427 | + raise | ||
| 428 | + except Exception as exc: | ||
| 429 | + raise RedisUnavailable("Redis write failed") from exc | ||
| 430 | + return { | ||
| 431 | + "key": request.key, | ||
| 432 | + "value": request.value, | ||
| 433 | + "ttl_seconds": ttl, | ||
| 434 | + "backend_mode": config.redis_mode, | ||
| 435 | + } | ||
| 436 | + | ||
| 437 | + | ||
| 438 | + "redis/read", | ||
| 439 | + request_model=RedisReadInput, | ||
| 440 | + response_model=RedisReadOutput, | ||
| 441 | + summary="Read a Redis value", | ||
| 442 | + description="Read one namespaced Redis key.", | ||
| 443 | + tags=["Redis"], | ||
| 444 | + ) | ||
| 445 | + async def read_redis( | ||
| 446 | + ctx: TypedAppContext[RedisReadInput], | ||
| 447 | + env: Envelope[RedisReadInput], | ||
| 448 | + ) -> dict[str, Any]: | ||
| 449 | + try: | ||
| 450 | + value = await ctx.kv.get(ctx.request.key) | ||
| 451 | + except RedisUnavailable: | ||
| 452 | + raise | ||
| 453 | + except Exception as exc: | ||
| 454 | + raise RedisUnavailable("Redis read failed") from exc | ||
| 455 | + return { | ||
| 456 | + "key": ctx.request.key, | ||
| 457 | + "found": value is not None, | ||
| 458 | + "value": value, | ||
| 459 | + "backend_mode": config.redis_mode, | ||
| 460 | + } | ||
| 461 | + | ||
| 462 | + | ||
| 463 | + "envelope/inspect", | ||
| 464 | + request_model=EnvelopeInspectInput, | ||
| 465 | + response_model=EnvelopeInspectOutput, | ||
| 466 | + summary="Inspect an Envelope", | ||
| 467 | + description="Return parsed Envelope, Metadata, and request context fields.", | ||
| 468 | + tags=["Envelope"], | ||
| 469 | + ) | ||
| 470 | + async def inspect_envelope( | ||
| 471 | + ctx: TypedAppContext[EnvelopeInspectInput], | ||
| 472 | + env: Envelope[EnvelopeInspectInput], | ||
| 473 | + ) -> dict[str, Any]: | ||
| 474 | + return { | ||
| 475 | + "envelope": env.to_dict(), | ||
| 476 | + "context": { | ||
| 477 | + "msg_type": ctx.msg_type, | ||
| 478 | + "request_id": ctx.request_id, | ||
| 479 | + "user_id": ctx.user_id, | ||
| 480 | + "chat_id": ctx.chat_id, | ||
| 481 | + "session_id": ctx.session_id, | ||
| 482 | + "bot_id": ctx.bot_id, | ||
| 483 | + "channel": ctx.channel, | ||
| 484 | + "trace_id": ctx.trace_id, | ||
| 485 | + "instance_id": ctx.instance_id, | ||
| 486 | + "replica_id": ctx.replica_id, | ||
| 487 | + }, | ||
| 488 | + "service": { | ||
| 489 | + "environment": config.environment, | ||
| 490 | + "service_label": config.service_label, | ||
| 491 | + "redis_mode": config.redis_mode, | ||
| 492 | + }, | ||
| 493 | + } | ||
| 494 | + | ||
| 495 | + | ||
| 496 | + "k8s/pod/read", | ||
| 497 | + request_model=PodNameInput, | ||
| 498 | + response_model=PodReadOutput, | ||
| 499 | + summary="Read a managed Kubernetes Pod", | ||
| 500 | + description="Read lifecycle fields for one Demo-managed Pod.", | ||
| 501 | + tags=["Kubernetes"], | ||
| 502 | + ) | ||
| 503 | + async def read_pod( | ||
| 504 | + ctx: TypedAppContext[PodNameInput], | ||
| 505 | + env: Envelope[PodNameInput], | ||
| 506 | + ) -> dict[str, Any]: | ||
| 507 | + pod = await ctx.kubernetes.get_pod(ctx.request.name) | ||
| 508 | + if pod is None: | ||
| 509 | + raise NotFoundError(f"pod {ctx.request.name!r} not found") | ||
| 510 | + return {"pod": asdict(pod)} | ||
| 511 | + | ||
| 512 | + | ||
| 513 | + "k8s/pod/create", | ||
| 514 | + request_model=PodNameInput, | ||
| 515 | + response_model=PodCreateOutput, | ||
| 516 | + summary="Create a managed Kubernetes Pod", | ||
| 517 | + description="Create one Pod from the configured restricted template.", | ||
| 518 | + tags=["Kubernetes"], | ||
| 519 | + ) | ||
| 520 | + async def create_pod( | ||
| 521 | + ctx: TypedAppContext[PodNameInput], | ||
| 522 | + env: Envelope[PodNameInput], | ||
| 523 | + ) -> dict[str, Any]: | ||
| 524 | + pod = await ctx.kubernetes.create_pod( | ||
| 525 | + PodCreateSpec( | ||
| 526 | + name=ctx.request.name, | ||
| 527 | + image=config.kubernetes_pod_image, | ||
| 528 | + ) | ||
| 529 | + ) | ||
| 530 | + return {"operation": "created", "pod": asdict(pod)} | ||
| 531 | + | ||
| 532 | + | ||
| 533 | + "k8s/pod/delete", | ||
| 534 | + request_model=PodNameInput, | ||
| 535 | + response_model=PodDeleteOutput, | ||
| 536 | + summary="Delete a managed Kubernetes Pod", | ||
| 537 | + description="Submit deletion for one Demo-managed Pod.", | ||
| 538 | + tags=["Kubernetes"], | ||
| 539 | + ) | ||
| 540 | + async def delete_pod( | ||
| 541 | + ctx: TypedAppContext[PodNameInput], | ||
| 542 | + env: Envelope[PodNameInput], | ||
| 543 | + ) -> dict[str, Any]: | ||
| 544 | + return asdict(await ctx.kubernetes.delete_pod(ctx.request.name)) | ||
| 545 | + | ||
| 546 | + | ||
| 547 | +def create_app(demo_config: DemoConfig | None = None) -> App: | ||
| 548 | + config = demo_config or DemoConfig.from_env() | ||
| 549 | + | ||
| 550 | + def create_system_context() -> SystemContext: | ||
| 551 | + service_config = ServiceConfig.from_env() | ||
| 552 | + return build_demo_system_context(service_config, config) | ||
| 553 | + | ||
| 554 | + application = App( | ||
| 555 | + create_system_context, | ||
| 556 | + title="Simple Service Capabilities Demo", | ||
| 557 | + ) | ||
| 558 | + register_handlers(application, config) | ||
| 559 | + return application | ||
| 560 | + | ||
| 561 | + | ||
| 562 | +app = create_app() | ||
| 563 | +asgi = app.asgi | ||
| 564 | + | ||
| 565 | + | ||
| 566 | +if __name__ == "__main__": | ||
| 567 | + app.run() | ||
| @@ -22,6 +22,8 @@ from .errors import ( | |||
| 22 | LockNotAcquired, | 22 | LockNotAcquired, |
| 23 | NotFoundError, | 23 | NotFoundError, |
| 24 | Interrupted, | 24 | Interrupted, |
| 25 | + KubernetesUnavailable, | ||
| 26 | + PermissionDenied, | ||
| 25 | RedisUnavailable, | 27 | RedisUnavailable, |
| 26 | ValidationError, | 28 | ValidationError, |
| 27 | ) | 29 | ) |
| @@ -44,12 +46,18 @@ from .context import ( | |||
| 44 | LockLease, | 46 | LockLease, |
| 45 | LockManager, | 47 | LockManager, |
| 46 | JsonCacheSerializer, | 48 | JsonCacheSerializer, |
| 49 | + FakeKubernetesOperations, | ||
| 50 | + KubernetesAsyncioOperations, | ||
| 51 | + KubernetesOperations, | ||
| 47 | MemoryCacheBackend, | 52 | MemoryCacheBackend, |
| 48 | MemoryLockBackend, | 53 | MemoryLockBackend, |
| 49 | NoopAuditLogger, | 54 | NoopAuditLogger, |
| 50 | RedisLockBackend, | 55 | RedisLockBackend, |
| 51 | RedisCacheBackend, | 56 | RedisCacheBackend, |
| 52 | RequestContext, | 57 | RequestContext, |
| 58 | + PodCreateSpec, | ||
| 59 | + PodDeleteResult, | ||
| 60 | + PodSummary, | ||
| 53 | SystemContext, | 61 | SystemContext, |
| 54 | TypedAppContext, | 62 | TypedAppContext, |
| 55 | build_lock_backend, | 63 | build_lock_backend, |
| @@ -97,6 +105,8 @@ __all__ = [ | |||
| 97 | "DatabaseUnavailable", | 105 | "DatabaseUnavailable", |
| 98 | "RedisUnavailable", | 106 | "RedisUnavailable", |
| 99 | "CacheUnavailable", | 107 | "CacheUnavailable", |
| 108 | + "KubernetesUnavailable", | ||
| 109 | + "PermissionDenied", | ||
| 100 | "ValidationError", | 110 | "ValidationError", |
| 101 | "NotFoundError", | 111 | "NotFoundError", |
| 102 | "IdempotentConflict", | 112 | "IdempotentConflict", |
| @@ -120,6 +130,12 @@ __all__ = [ | |||
| 120 | "EtcdLockBackend", | 130 | "EtcdLockBackend", |
| 121 | "LoggingAuditLogger", | 131 | "LoggingAuditLogger", |
| 122 | "NoopAuditLogger", | 132 | "NoopAuditLogger", |
| 133 | + "FakeKubernetesOperations", | ||
| 134 | + "KubernetesAsyncioOperations", | ||
| 135 | + "KubernetesOperations", | ||
| 136 | + "PodCreateSpec", | ||
| 137 | + "PodDeleteResult", | ||
| 138 | + "PodSummary", | ||
| 123 | # locks | 139 | # locks |
| 124 | "LeaseState", | 140 | "LeaseState", |
| 125 | "LockBackend", | 141 | "LockBackend", |
| @@ -162,6 +162,7 @@ def build_system_context( | |||
| 162 | *, | 162 | *, |
| 163 | db: Any = None, | 163 | db: Any = None, |
| 164 | redis: Any = None, | 164 | redis: Any = None, |
| 165 | + kubernetes: Any = None, | ||
| 165 | etcd_client: Any = None, | 166 | etcd_client: Any = None, |
| 166 | lock_backend: Any = None, | 167 | lock_backend: Any = None, |
| 167 | cache_backend: Any = None, | 168 | cache_backend: Any = None, |
| @@ -191,6 +192,7 @@ def build_system_context( | |||
| 191 | cache_resource = build_cache_backend(cfg, redis_client=redis_resource) | 192 | cache_resource = build_cache_backend(cfg, redis_client=redis_resource) |
| 192 | return SystemContext( | 193 | return SystemContext( |
| 193 | redis=redis_resource, | 194 | redis=redis_resource, |
| 195 | + kubernetes=kubernetes, | ||
| 194 | db=db_resource, | 196 | db=db_resource, |
| 195 | settings=cfg, | 197 | settings=cfg, |
| 196 | key_prefix=cfg.key_prefix, | 198 | key_prefix=cfg.key_prefix, |
| @@ -29,6 +29,14 @@ from .locks import ( | |||
| 29 | build_lock_backend, | 29 | build_lock_backend, |
| 30 | create_lock_backend, | 30 | create_lock_backend, |
| 31 | ) | 31 | ) |
| 32 | +from .kubernetes import ( | ||
| 33 | + FakeKubernetesOperations, | ||
| 34 | + KubernetesAsyncioOperations, | ||
| 35 | + KubernetesOperations, | ||
| 36 | + PodCreateSpec, | ||
| 37 | + PodDeleteResult, | ||
| 38 | + PodSummary, | ||
| 39 | +) | ||
| 32 | from .request_context import RequestContext, TypedAppContext | 40 | from .request_context import RequestContext, TypedAppContext |
| 33 | from .system_context import SystemContext | 41 | from .system_context import SystemContext |
| 34 | 42 | ||
| @@ -59,6 +67,12 @@ __all__ = [ | |||
| 59 | "SystemContext", | 67 | "SystemContext", |
| 60 | "TypedAppContext", | 68 | "TypedAppContext", |
| 61 | "JsonCacheSerializer", | 69 | "JsonCacheSerializer", |
| 70 | + "FakeKubernetesOperations", | ||
| 71 | + "KubernetesAsyncioOperations", | ||
| 72 | + "KubernetesOperations", | ||
| 73 | + "PodCreateSpec", | ||
| 74 | + "PodDeleteResult", | ||
| 75 | + "PodSummary", | ||
| 62 | "build_cache_backend", | 76 | "build_cache_backend", |
| 63 | "build_lock_backend", | 77 | "build_lock_backend", |
| 64 | "create_cache_backend", | 78 | "create_cache_backend", |
| @@ -0,0 +1,20 @@ | |||
| 1 | +# coding: utf-8 | ||
| 2 | +# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved | ||
| 3 | + | ||
| 4 | +from .asyncio_client import KubernetesAsyncioOperations | ||
| 5 | +from .base import ( | ||
| 6 | + KubernetesOperations, | ||
| 7 | + PodCreateSpec, | ||
| 8 | + PodDeleteResult, | ||
| 9 | + PodSummary, | ||
| 10 | +) | ||
| 11 | +from .fake import FakeKubernetesOperations | ||
| 12 | + | ||
| 13 | +__all__ = [ | ||
| 14 | + "FakeKubernetesOperations", | ||
| 15 | + "KubernetesAsyncioOperations", | ||
| 16 | + "KubernetesOperations", | ||
| 17 | + "PodCreateSpec", | ||
| 18 | + "PodDeleteResult", | ||
| 19 | + "PodSummary", | ||
| 20 | +] | ||
| @@ -0,0 +1,302 @@ | |||
| 1 | +# coding: utf-8 | ||
| 2 | +# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved | ||
| 3 | + | ||
| 4 | +"""Kubernetes Pod operations backed by ``kubernetes_asyncio``.""" | ||
| 5 | + | ||
| 6 | +from __future__ import annotations | ||
| 7 | + | ||
| 8 | +from inspect import isawaitable | ||
| 9 | +from pathlib import Path | ||
| 10 | +from typing import Any, Literal, Mapping | ||
| 11 | + | ||
| 12 | +from ...errors import ( | ||
| 13 | + ErrorCode, | ||
| 14 | + FrameworkError, | ||
| 15 | + KubernetesUnavailable, | ||
| 16 | + NotFoundError, | ||
| 17 | + PermissionDenied, | ||
| 18 | +) | ||
| 19 | +from .base import PodCreateSpec, PodDeleteResult, PodSummary | ||
| 20 | + | ||
| 21 | + | ||
| 22 | +class KubernetesAsyncioOperations: | ||
| 23 | + """Operate managed Pods in one namespace through a shared API client.""" | ||
| 24 | + | ||
| 25 | + REQUEST_TIMEOUT_SECONDS = 10 | ||
| 26 | + DELETE_GRACE_PERIOD_SECONDS = 0 | ||
| 27 | + DEFAULT_LABELS = { | ||
| 28 | + "app.kubernetes.io/managed-by": "openjiuwen-service-demo", | ||
| 29 | + } | ||
| 30 | + | ||
| 31 | + def __init__( | ||
| 32 | + self, | ||
| 33 | + namespace: str, | ||
| 34 | + *, | ||
| 35 | + labels: Mapping[str, str] | None = None, | ||
| 36 | + kubeconfig: str | Path | None = None, | ||
| 37 | + ) -> None: | ||
| 38 | + self.namespace = namespace | ||
| 39 | + self.labels = dict(labels or self.DEFAULT_LABELS) | ||
| 40 | + self.kubeconfig = str(kubeconfig) if kubeconfig else None | ||
| 41 | + self._api_client: Any = None | ||
| 42 | + self._core_api: Any = None | ||
| 43 | + self._client_module: Any = None | ||
| 44 | + | ||
| 45 | + async def start(self) -> None: | ||
| 46 | + if self._core_api is not None: | ||
| 47 | + return | ||
| 48 | + try: | ||
| 49 | + from kubernetes_asyncio import client, config | ||
| 50 | + from kubernetes_asyncio.config.config_exception import ConfigException | ||
| 51 | + except Exception as exc: | ||
| 52 | + raise KubernetesUnavailable( | ||
| 53 | + "kubernetes_asyncio is required for real Kubernetes mode" | ||
| 54 | + ) from exc | ||
| 55 | + | ||
| 56 | + api_client = None | ||
| 57 | + try: | ||
| 58 | + try: | ||
| 59 | + result = config.load_incluster_config() | ||
| 60 | + if isawaitable(result): | ||
| 61 | + await result | ||
| 62 | + except ConfigException: | ||
| 63 | + result = config.load_kube_config(config_file=self.kubeconfig) | ||
| 64 | + if isawaitable(result): | ||
| 65 | + await result | ||
| 66 | + api_client = client.ApiClient() | ||
| 67 | + self._client_module = client | ||
| 68 | + self._api_client = api_client | ||
| 69 | + self._core_api = client.CoreV1Api(api_client) | ||
| 70 | + except Exception as exc: | ||
| 71 | + if api_client is not None: | ||
| 72 | + result = api_client.close() | ||
| 73 | + if isawaitable(result): | ||
| 74 | + await result | ||
| 75 | + self._client_module = None | ||
| 76 | + self._api_client = None | ||
| 77 | + self._core_api = None | ||
| 78 | + raise KubernetesUnavailable( | ||
| 79 | + f"cannot initialize Kubernetes client: {exc}" | ||
| 80 | + ) from exc | ||
| 81 | + | ||
| 82 | + async def ping(self) -> bool: | ||
| 83 | + api = self._require_started() | ||
| 84 | + try: | ||
| 85 | + await api.list_namespaced_pod( | ||
| 86 | + namespace=self.namespace, | ||
| 87 | + limit=1, | ||
| 88 | + _request_timeout=self.REQUEST_TIMEOUT_SECONDS, | ||
| 89 | + ) | ||
| 90 | + return True | ||
| 91 | + except Exception as exc: | ||
| 92 | + raise self._map_api_exception(exc, operation="list Pods") from exc | ||
| 93 | + | ||
| 94 | + async def close(self) -> None: | ||
| 95 | + api_client = self._api_client | ||
| 96 | + self._core_api = None | ||
| 97 | + self._api_client = None | ||
| 98 | + self._client_module = None | ||
| 99 | + if api_client is not None: | ||
| 100 | + result = api_client.close() | ||
| 101 | + if isawaitable(result): | ||
| 102 | + await result | ||
| 103 | + | ||
| 104 | + async def get_pod(self, name: str) -> PodSummary | None: | ||
| 105 | + api = self._require_started() | ||
| 106 | + try: | ||
| 107 | + pod = await api.read_namespaced_pod( | ||
| 108 | + name=name, | ||
| 109 | + namespace=self.namespace, | ||
| 110 | + _request_timeout=self.REQUEST_TIMEOUT_SECONDS, | ||
| 111 | + ) | ||
| 112 | + except Exception as exc: | ||
| 113 | + if self._status(exc) == 404: | ||
| 114 | + return None | ||
| 115 | + raise self._map_api_exception(exc, operation=f"read Pod {name!r}") from exc | ||
| 116 | + if not self._is_managed(pod): | ||
| 117 | + return None | ||
| 118 | + return self._to_summary(pod) | ||
| 119 | + | ||
| 120 | + async def create_pod(self, spec: PodCreateSpec) -> PodSummary: | ||
| 121 | + api = self._require_started() | ||
| 122 | + client = self._client_module | ||
| 123 | + container = client.V1Container( | ||
| 124 | + name="demo", | ||
| 125 | + image=spec.image, | ||
| 126 | + image_pull_policy="IfNotPresent", | ||
| 127 | + security_context=client.V1SecurityContext( | ||
| 128 | + allow_privilege_escalation=False, | ||
| 129 | + read_only_root_filesystem=False, | ||
| 130 | + run_as_non_root=True, | ||
| 131 | + seccomp_profile=client.V1SeccompProfile(type="RuntimeDefault"), | ||
| 132 | + capabilities=client.V1Capabilities(drop=["ALL"]), | ||
| 133 | + ), | ||
| 134 | + resources=client.V1ResourceRequirements( | ||
| 135 | + requests={"cpu": "10m", "memory": "16Mi"}, | ||
| 136 | + limits={"cpu": "100m", "memory": "64Mi"}, | ||
| 137 | + ), | ||
| 138 | + ) | ||
| 139 | + body = client.V1Pod( | ||
| 140 | + api_version="v1", | ||
| 141 | + kind="Pod", | ||
| 142 | + metadata=client.V1ObjectMeta( | ||
| 143 | + name=spec.name, | ||
| 144 | + namespace=self.namespace, | ||
| 145 | + labels=dict(self.labels), | ||
| 146 | + ), | ||
| 147 | + spec=client.V1PodSpec( | ||
| 148 | + automount_service_account_token=False, | ||
| 149 | + containers=[container], | ||
| 150 | + restart_policy="Never", | ||
| 151 | + ), | ||
| 152 | + ) | ||
| 153 | + try: | ||
| 154 | + pod = await api.create_namespaced_pod( | ||
| 155 | + namespace=self.namespace, | ||
| 156 | + body=body, | ||
| 157 | + _request_timeout=self.REQUEST_TIMEOUT_SECONDS, | ||
| 158 | + ) | ||
| 159 | + except Exception as exc: | ||
| 160 | + raise self._map_api_exception( | ||
| 161 | + exc, | ||
| 162 | + operation=f"create Pod {spec.name!r}", | ||
| 163 | + ) from exc | ||
| 164 | + return self._to_summary(pod) | ||
| 165 | + | ||
| 166 | + async def delete_pod(self, name: str) -> PodDeleteResult: | ||
| 167 | + api = self._require_started() | ||
| 168 | + try: | ||
| 169 | + pod = await api.read_namespaced_pod( | ||
| 170 | + name=name, | ||
| 171 | + namespace=self.namespace, | ||
| 172 | + _request_timeout=self.REQUEST_TIMEOUT_SECONDS, | ||
| 173 | + ) | ||
| 174 | + except Exception as exc: | ||
| 175 | + if self._status(exc) == 404: | ||
| 176 | + return self._delete_result(name, "already_absent") | ||
| 177 | + raise self._map_api_exception(exc, operation=f"read Pod {name!r}") from exc | ||
| 178 | + | ||
| 179 | + if not self._is_managed(pod): | ||
| 180 | + raise NotFoundError(f"pod {name!r} not found") | ||
| 181 | + metadata = self._value(pod, "metadata") | ||
| 182 | + if self._value(metadata, "deletion_timestamp") is not None: | ||
| 183 | + return self._delete_result(name, "deletion_in_progress") | ||
| 184 | + uid = self._value(metadata, "uid") | ||
| 185 | + if not uid: | ||
| 186 | + raise FrameworkError(f"Pod {name!r} response has no UID") | ||
| 187 | + | ||
| 188 | + client = self._client_module | ||
| 189 | + body = client.V1DeleteOptions( | ||
| 190 | + grace_period_seconds=self.DELETE_GRACE_PERIOD_SECONDS, | ||
| 191 | + preconditions=client.V1Preconditions(uid=uid), | ||
| 192 | + ) | ||
| 193 | + try: | ||
| 194 | + await api.delete_namespaced_pod( | ||
| 195 | + name=name, | ||
| 196 | + namespace=self.namespace, | ||
| 197 | + body=body, | ||
| 198 | + grace_period_seconds=self.DELETE_GRACE_PERIOD_SECONDS, | ||
| 199 | + _request_timeout=self.REQUEST_TIMEOUT_SECONDS, | ||
| 200 | + ) | ||
| 201 | + except Exception as exc: | ||
| 202 | + if self._status(exc) == 404: | ||
| 203 | + return self._delete_result(name, "already_absent") | ||
| 204 | + raise self._map_api_exception( | ||
| 205 | + exc, | ||
| 206 | + operation=f"delete Pod {name!r}", | ||
| 207 | + ) from exc | ||
| 208 | + return self._delete_result(name, "delete_requested") | ||
| 209 | + | ||
| 210 | + def _require_started(self) -> Any: | ||
| 211 | + if self._core_api is None: | ||
| 212 | + raise KubernetesUnavailable("Kubernetes operations are not started") | ||
| 213 | + return self._core_api | ||
| 214 | + | ||
| 215 | + def _is_managed(self, pod: Any) -> bool: | ||
| 216 | + metadata = self._value(pod, "metadata") | ||
| 217 | + labels = self._value(metadata, "labels") or {} | ||
| 218 | + return all(labels.get(key) == value for key, value in self.labels.items()) | ||
| 219 | + | ||
| 220 | + def _to_summary(self, pod: Any) -> PodSummary: | ||
| 221 | + metadata = self._value(pod, "metadata") | ||
| 222 | + spec = self._value(pod, "spec") | ||
| 223 | + status = self._value(pod, "status") | ||
| 224 | + name = self._value(metadata, "name") | ||
| 225 | + namespace = self._value(metadata, "namespace") or self.namespace | ||
| 226 | + if not name: | ||
| 227 | + raise FrameworkError("Kubernetes Pod response has no name") | ||
| 228 | + | ||
| 229 | + phase = self._value(status, "phase") or "Unknown" | ||
| 230 | + if self._value(metadata, "deletion_timestamp") is not None: | ||
| 231 | + phase = "Terminating" | ||
| 232 | + conditions = self._value(status, "conditions") or [] | ||
| 233 | + ready_condition = any( | ||
| 234 | + self._value(condition, "type") == "Ready" | ||
| 235 | + and str(self._value(condition, "status")).lower() == "true" | ||
| 236 | + for condition in conditions | ||
| 237 | + ) | ||
| 238 | + container_statuses = self._value(status, "container_statuses") or [] | ||
| 239 | + containers_ready = bool(container_statuses) and all( | ||
| 240 | + bool(self._value(item, "ready")) for item in container_statuses | ||
| 241 | + ) | ||
| 242 | + containers = self._value(spec, "containers") or [] | ||
| 243 | + image = self._value(containers[0], "image") if containers else None | ||
| 244 | + return PodSummary( | ||
| 245 | + name=str(name), | ||
| 246 | + namespace=str(namespace), | ||
| 247 | + phase=str(phase), | ||
| 248 | + ready=ready_condition and containers_ready, | ||
| 249 | + image=str(image) if image is not None else None, | ||
| 250 | + ) | ||
| 251 | + | ||
| 252 | + def _delete_result( | ||
| 253 | + self, | ||
| 254 | + name: str, | ||
| 255 | + state: Literal[ | ||
| 256 | + "delete_requested", | ||
| 257 | + "deletion_in_progress", | ||
| 258 | + "already_absent", | ||
| 259 | + ], | ||
| 260 | + ) -> PodDeleteResult: | ||
| 261 | + return PodDeleteResult( | ||
| 262 | + name=name, | ||
| 263 | + namespace=self.namespace, | ||
| 264 | + state=state, | ||
| 265 | + ) | ||
| 266 | + | ||
| 267 | + | ||
| 268 | + def _map_api_exception(cls, exc: Exception, *, operation: str) -> FrameworkError: | ||
| 269 | + status = cls._status(exc) | ||
| 270 | + reason = getattr(exc, "reason", None) or str(exc) or exc.__class__.__name__ | ||
| 271 | + message = f"Kubernetes {operation} failed" | ||
| 272 | + if status is not None: | ||
| 273 | + message += f" with status {status}" | ||
| 274 | + message += f": {reason}" | ||
| 275 | + if status in {401, 403}: | ||
| 276 | + return PermissionDenied(message) | ||
| 277 | + if status == 409: | ||
| 278 | + return FrameworkError(message, code=ErrorCode.CONFLICT) | ||
| 279 | + if status in {429, 500, 502, 503, 504} or status is None: | ||
| 280 | + return KubernetesUnavailable(message) | ||
| 281 | + if status == 404: | ||
| 282 | + return NotFoundError(message) | ||
| 283 | + return FrameworkError(message) | ||
| 284 | + | ||
| 285 | + | ||
| 286 | + def _status(exc: Exception) -> int | None: | ||
| 287 | + value = getattr(exc, "status", None) | ||
| 288 | + try: | ||
| 289 | + return int(value) if value is not None else None | ||
| 290 | + except (TypeError, ValueError): | ||
| 291 | + return None | ||
| 292 | + | ||
| 293 | + | ||
| 294 | + def _value(obj: Any, name: str) -> Any: | ||
| 295 | + if obj is None: | ||
| 296 | + return None | ||
| 297 | + if isinstance(obj, Mapping): | ||
| 298 | + return obj.get(name) | ||
| 299 | + return getattr(obj, name, None) | ||
| 300 | + | ||
| 301 | + | ||
| 302 | +__all__ = ["KubernetesAsyncioOperations"] | ||
| @@ -0,0 +1,72 @@ | |||
| 1 | +# coding: utf-8 | ||
| 2 | +# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved | ||
| 3 | + | ||
| 4 | +"""Kubernetes Pod capability contracts and service-domain models.""" | ||
| 5 | + | ||
| 6 | +from __future__ import annotations | ||
| 7 | + | ||
| 8 | +from dataclasses import dataclass | ||
| 9 | +from typing import Literal, Protocol, runtime_checkable | ||
| 10 | + | ||
| 11 | + | ||
| 12 | + | ||
| 13 | +class PodCreateSpec: | ||
| 14 | + """Restricted input used to create one managed Pod.""" | ||
| 15 | + | ||
| 16 | + name: str | ||
| 17 | + image: str | ||
| 18 | + | ||
| 19 | + | ||
| 20 | + | ||
| 21 | +class PodSummary: | ||
| 22 | + """Stable Pod lifecycle fields exposed to service handlers.""" | ||
| 23 | + | ||
| 24 | + name: str | ||
| 25 | + namespace: str | ||
| 26 | + phase: str | ||
| 27 | + ready: bool | ||
| 28 | + image: str | None | ||
| 29 | + | ||
| 30 | + | ||
| 31 | + | ||
| 32 | +class PodDeleteResult: | ||
| 33 | + """Result of submitting or observing a Pod deletion.""" | ||
| 34 | + | ||
| 35 | + name: str | ||
| 36 | + namespace: str | ||
| 37 | + state: Literal[ | ||
| 38 | + "delete_requested", | ||
| 39 | + "deletion_in_progress", | ||
| 40 | + "already_absent", | ||
| 41 | + ] | ||
| 42 | + | ||
| 43 | + | ||
| 44 | + | ||
| 45 | +class KubernetesOperations(Protocol): | ||
| 46 | + """Process-level asynchronous operations for managed Pods.""" | ||
| 47 | + | ||
| 48 | + async def start(self) -> None: | ||
| 49 | + ... | ||
| 50 | + | ||
| 51 | + async def ping(self) -> bool: | ||
| 52 | + ... | ||
| 53 | + | ||
| 54 | + async def close(self) -> None: | ||
| 55 | + ... | ||
| 56 | + | ||
| 57 | + async def get_pod(self, name: str) -> PodSummary | None: | ||
| 58 | + ... | ||
| 59 | + | ||
| 60 | + async def create_pod(self, spec: PodCreateSpec) -> PodSummary: | ||
| 61 | + ... | ||
| 62 | + | ||
| 63 | + async def delete_pod(self, name: str) -> PodDeleteResult: | ||
| 64 | + ... | ||
| 65 | + | ||
| 66 | + | ||
| 67 | +__all__ = [ | ||
| 68 | + "KubernetesOperations", | ||
| 69 | + "PodCreateSpec", | ||
| 70 | + "PodDeleteResult", | ||
| 71 | + "PodSummary", | ||
| 72 | +] | ||
| @@ -0,0 +1,86 @@ | |||
| 1 | +# coding: utf-8 | ||
| 2 | +# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved | ||
| 3 | + | ||
| 4 | +"""In-memory Kubernetes Pod operations for local development and tests.""" | ||
| 5 | + | ||
| 6 | +from __future__ import annotations | ||
| 7 | + | ||
| 8 | +import asyncio | ||
| 9 | +from dataclasses import replace | ||
| 10 | +from typing import Mapping | ||
| 11 | + | ||
| 12 | +from ...errors import ErrorCode, FrameworkError, KubernetesUnavailable | ||
| 13 | +from .base import PodCreateSpec, PodDeleteResult, PodSummary | ||
| 14 | + | ||
| 15 | + | ||
| 16 | +class FakeKubernetesOperations: | ||
| 17 | + """Keep managed Pod state in one service process.""" | ||
| 18 | + | ||
| 19 | + def __init__( | ||
| 20 | + self, | ||
| 21 | + namespace: str, | ||
| 22 | + *, | ||
| 23 | + labels: Mapping[str, str] | None = None, | ||
| 24 | + ) -> None: | ||
| 25 | + self.namespace = namespace | ||
| 26 | + self.labels = dict(labels or {}) | ||
| 27 | + self._pods: dict[str, PodSummary] = {} | ||
| 28 | + self._lock: asyncio.Lock | None = None | ||
| 29 | + self._started = False | ||
| 30 | + | ||
| 31 | + async def start(self) -> None: | ||
| 32 | + if self._started: | ||
| 33 | + return | ||
| 34 | + self._lock = asyncio.Lock() | ||
| 35 | + self._started = True | ||
| 36 | + | ||
| 37 | + async def ping(self) -> bool: | ||
| 38 | + return self._started | ||
| 39 | + | ||
| 40 | + async def close(self) -> None: | ||
| 41 | + self._started = False | ||
| 42 | + self._pods.clear() | ||
| 43 | + self._lock = None | ||
| 44 | + | ||
| 45 | + async def get_pod(self, name: str) -> PodSummary | None: | ||
| 46 | + lock = self._require_started() | ||
| 47 | + async with lock: | ||
| 48 | + pod = self._pods.get(name) | ||
| 49 | + return replace(pod) if pod is not None else None | ||
| 50 | + | ||
| 51 | + async def create_pod(self, spec: PodCreateSpec) -> PodSummary: | ||
| 52 | + lock = self._require_started() | ||
| 53 | + async with lock: | ||
| 54 | + if spec.name in self._pods: | ||
| 55 | + raise FrameworkError( | ||
| 56 | + f"pod {spec.name!r} already exists", | ||
| 57 | + code=ErrorCode.CONFLICT, | ||
| 58 | + ) | ||
| 59 | + pod = PodSummary( | ||
| 60 | + name=spec.name, | ||
| 61 | + namespace=self.namespace, | ||
| 62 | + phase="Running", | ||
| 63 | + ready=True, | ||
| 64 | + image=spec.image, | ||
| 65 | + ) | ||
| 66 | + self._pods[spec.name] = pod | ||
| 67 | + return replace(pod) | ||
| 68 | + | ||
| 69 | + async def delete_pod(self, name: str) -> PodDeleteResult: | ||
| 70 | + lock = self._require_started() | ||
| 71 | + async with lock: | ||
| 72 | + pod = self._pods.pop(name, None) | ||
| 73 | + state = "delete_requested" if pod is not None else "already_absent" | ||
| 74 | + return PodDeleteResult( | ||
| 75 | + name=name, | ||
| 76 | + namespace=self.namespace, | ||
| 77 | + state=state, | ||
| 78 | + ) | ||
| 79 | + | ||
| 80 | + def _require_started(self) -> asyncio.Lock: | ||
| 81 | + if not self._started or self._lock is None: | ||
| 82 | + raise KubernetesUnavailable("Kubernetes operations are not started") | ||
| 83 | + return self._lock | ||
| 84 | + | ||
| 85 | + | ||
| 86 | +__all__ = ["FakeKubernetesOperations"] | ||
| @@ -24,6 +24,7 @@ from .primitives.kv_store import KVStore | |||
| 24 | 24 | ||
| 25 | if TYPE_CHECKING: | 25 | if TYPE_CHECKING: |
| 26 | from .cache import Cache | 26 | from .cache import Cache |
| 27 | + from .kubernetes import KubernetesOperations | ||
| 27 | from .locks import LockManager | 28 | from .locks import LockManager |
| 28 | from .system_context import SystemContext | 29 | from .system_context import SystemContext |
| 29 | 30 | ||
| @@ -260,6 +261,11 @@ class RequestContext(Generic[TRequest]): | |||
| 260 | """ | 261 | """ |
| 261 | return self.require_redis() | 262 | return self.require_redis() |
| 262 | 263 | ||
| 264 | + | ||
| 265 | + def kubernetes(self) -> KubernetesOperations: | ||
| 266 | + """Return the shared Kubernetes operations for this request.""" | ||
| 267 | + return self.require_kubernetes() | ||
| 268 | + | ||
| 263 | def require_redis(self) -> Any: | 269 | def require_redis(self) -> Any: |
| 264 | """Require the shared asynchronous Redis client for an active request.""" | 270 | """Require the shared asynchronous Redis client for an active request.""" |
| 265 | self.check_interrupted() | 271 | self.check_interrupted() |
| @@ -270,6 +276,11 @@ class RequestContext(Generic[TRequest]): | |||
| 270 | self.check_interrupted() | 276 | self.check_interrupted() |
| 271 | return self.sysctx.require_db() | 277 | return self.sysctx.require_db() |
| 272 | 278 | ||
| 279 | + def require_kubernetes(self) -> KubernetesOperations: | ||
| 280 | + """Require shared Kubernetes operations for an active request.""" | ||
| 281 | + self.check_interrupted() | ||
| 282 | + return self.sysctx.require_kubernetes() | ||
| 283 | + | ||
| 273 | async def db_create(self, table_name: str, data: dict[str, Any]) -> Any: | 284 | async def db_create(self, table_name: str, data: dict[str, Any]) -> Any: |
| 274 | """Create a record using an independent DBHandler operation.""" | 285 | """Create a record using an independent DBHandler operation.""" |
| 275 | return await self.require_db().create(table_name, data) | 286 | return await self.require_db().create(table_name, data) |
| @@ -31,8 +31,10 @@ from ..errors import ( | |||
| 31 | DatabaseUnavailable, | 31 | DatabaseUnavailable, |
| 32 | FrameworkError, | 32 | FrameworkError, |
| 33 | LockBackendUnavailable, | 33 | LockBackendUnavailable, |
| 34 | + KubernetesUnavailable, | ||
| 34 | RedisUnavailable, | 35 | RedisUnavailable, |
| 35 | ) | 36 | ) |
| 37 | +from .kubernetes import KubernetesOperations | ||
| 36 | from .audit import AuditEvent, AuditLogger, LoggingAuditLogger, NoopAuditLogger | 38 | from .audit import AuditEvent, AuditLogger, LoggingAuditLogger, NoopAuditLogger |
| 37 | from .request_context import RequestContext | 39 | from .request_context import RequestContext |
| 38 | 40 | ||
| @@ -58,12 +60,14 @@ class SystemContext: | |||
| 58 | lock_backend: Any = None, | 60 | lock_backend: Any = None, |
| 59 | cache_backend: Any = None, | 61 | cache_backend: Any = None, |
| 60 | cache: Any = None, | 62 | cache: Any = None, |
| 63 | + kubernetes: KubernetesOperations | None = None, | ||
| 61 | table_definitions: Iterable[Any] | None = None, | 64 | table_definitions: Iterable[Any] | None = None, |
| 62 | request_timeout_seconds: float | None = None, | 65 | request_timeout_seconds: float | None = None, |
| 63 | _owns_db: bool | None = None, | 66 | _owns_db: bool | None = None, |
| 64 | _owns_redis: bool = False, | 67 | _owns_redis: bool = False, |
| 65 | _owns_lock_backend: bool | None = None, | 68 | _owns_lock_backend: bool | None = None, |
| 66 | _owns_cache_backend: bool | None = None, | 69 | _owns_cache_backend: bool | None = None, |
| 70 | + _owns_kubernetes: bool = False, | ||
| 67 | ) -> None: | 71 | ) -> None: |
| 68 | self.redis = redis | 72 | self.redis = redis |
| 69 | self.db = db | 73 | self.db = db |
| @@ -88,6 +92,7 @@ class SystemContext: | |||
| 88 | if cache_backend is not None and cache is not None: | 92 | if cache_backend is not None and cache is not None: |
| 89 | raise ValueError("cache_backend and cache cannot both be provided") | 93 | raise ValueError("cache_backend and cache cannot both be provided") |
| 90 | self.cache_backend = cache_backend if cache_backend is not None else cache | 94 | self.cache_backend = cache_backend if cache_backend is not None else cache |
| 95 | + self.kubernetes = kubernetes | ||
| 91 | self.table_definitions = tuple(table_definitions or ()) | 96 | self.table_definitions = tuple(table_definitions or ()) |
| 92 | self._owns_db = db is not None if _owns_db is None else bool(_owns_db) | 97 | self._owns_db = db is not None if _owns_db is None else bool(_owns_db) |
| 93 | self._owns_redis = _owns_redis | 98 | self._owns_redis = _owns_redis |
| @@ -102,6 +107,7 @@ class SystemContext: | |||
| 102 | if _owns_cache_backend is None | 107 | if _owns_cache_backend is None |
| 103 | else bool(_owns_cache_backend) | 108 | else bool(_owns_cache_backend) |
| 104 | ) | 109 | ) |
| 110 | + self._owns_kubernetes = bool(kubernetes is not None and _owns_kubernetes) | ||
| 105 | self._started = False | 111 | self._started = False |
| 106 | self._stopped = False | 112 | self._stopped = False |
| 107 | self._active_resources: list[str] = [] | 113 | self._active_resources: list[str] = [] |
| @@ -141,6 +147,12 @@ class SystemContext: | |||
| 141 | raise CacheUnavailable("cache backend is not configured") | 147 | raise CacheUnavailable("cache backend is not configured") |
| 142 | return self.cache_backend | 148 | return self.cache_backend |
| 143 | 149 | ||
| 150 | + def require_kubernetes(self) -> KubernetesOperations: | ||
| 151 | + """Return configured Kubernetes operations or raise a framework error.""" | ||
| 152 | + if self.kubernetes is None: | ||
| 153 | + raise KubernetesUnavailable("Kubernetes operations are not configured") | ||
| 154 | + return self.kubernetes | ||
| 155 | + | ||
| 144 | def set_db(self, db: Any, *, owned: bool = False) -> None: | 156 | def set_db(self, db: Any, *, owned: bool = False) -> None: |
| 145 | self.db = db | 157 | self.db = db |
| 146 | self._owns_db = bool(db is not None and owned) | 158 | self._owns_db = bool(db is not None and owned) |
| @@ -159,6 +171,15 @@ class SystemContext: | |||
| 159 | self.cache_backend = backend | 171 | self.cache_backend = backend |
| 160 | self._owns_cache_backend = bool(backend is not None and owned) | 172 | self._owns_cache_backend = bool(backend is not None and owned) |
| 161 | 173 | ||
| 174 | + def set_kubernetes( | ||
| 175 | + self, | ||
| 176 | + kubernetes: KubernetesOperations | None, | ||
| 177 | + *, | ||
| 178 | + owned: bool = False, | ||
| 179 | + ) -> None: | ||
| 180 | + self.kubernetes = kubernetes | ||
| 181 | + self._owns_kubernetes = bool(kubernetes is not None and owned) | ||
| 182 | + | ||
| 162 | def set_audit_logger(self, audit_logger: AuditLogger | None) -> None: | 183 | def set_audit_logger(self, audit_logger: AuditLogger | None) -> None: |
| 163 | """Replace the process audit sink; ``None`` selects a no-op sink.""" | 184 | """Replace the process audit sink; ``None`` selects a no-op sink.""" |
| 164 | self.audit_logger = audit_logger or NoopAuditLogger() | 185 | self.audit_logger = audit_logger or NoopAuditLogger() |
| @@ -197,6 +218,15 @@ class SystemContext: | |||
| 197 | if not await self._ping(self.redis): | 218 | if not await self._ping(self.redis): |
| 198 | raise RedisUnavailable("Redis readiness check failed") | 219 | raise RedisUnavailable("Redis readiness check failed") |
| 199 | 220 | ||
| 221 | + if self.kubernetes is not None: | ||
| 222 | + self._active_resources.append("kubernetes") | ||
| 223 | + if self._owns_kubernetes: | ||
| 224 | + await self._call(self.kubernetes.start) | ||
| 225 | + if not await self._ping(self.kubernetes): | ||
| 226 | + raise KubernetesUnavailable( | ||
| 227 | + "Kubernetes readiness check failed" | ||
| 228 | + ) | ||
| 229 | + | ||
| 200 | if self.lock_backend is not None: | 230 | if self.lock_backend is not None: |
| 201 | self._active_resources.append("lock") | 231 | self._active_resources.append("lock") |
| 202 | connect = getattr(self.lock_backend, "connect", None) | 232 | connect = getattr(self.lock_backend, "connect", None) |
| @@ -223,6 +253,7 @@ class SystemContext: | |||
| 223 | resources = ( | 253 | resources = ( |
| 224 | ("db", self.db), | 254 | ("db", self.db), |
| 225 | ("redis", self.redis), | 255 | ("redis", self.redis), |
| 256 | + ("kubernetes", self.kubernetes), | ||
| 226 | ("lock", self.lock_backend), | 257 | ("lock", self.lock_backend), |
| 227 | ("cache", self.cache_backend), | 258 | ("cache", self.cache_backend), |
| 228 | ) | 259 | ) |
| @@ -238,12 +269,18 @@ class SystemContext: | |||
| 238 | statuses: dict[str, bool | None] = { | 269 | statuses: dict[str, bool | None] = { |
| 239 | "db": None, | 270 | "db": None, |
| 240 | "redis": None, | 271 | "redis": None, |
| 272 | + "kubernetes": None, | ||
| 241 | "lock": None, | 273 | "lock": None, |
| 242 | "cache": None, | 274 | "cache": None, |
| 243 | } | 275 | } |
| 244 | checks = ( | 276 | checks = ( |
| 245 | ("db", self.db, self._db_ready), | 277 | ("db", self.db, self._db_ready), |
| 246 | ("redis", self.redis, lambda: self._ping(self.redis)), | 278 | ("redis", self.redis, lambda: self._ping(self.redis)), |
| 279 | + ( | ||
| 280 | + "kubernetes", | ||
| 281 | + self.kubernetes, | ||
| 282 | + lambda: self._ping(self.kubernetes), | ||
| 283 | + ), | ||
| 247 | ("lock", self.lock_backend, lambda: self._ping(self.lock_backend)), | 284 | ("lock", self.lock_backend, lambda: self._ping(self.lock_backend)), |
| 248 | ("cache", self.cache_backend, lambda: self._ping(self.cache_backend)), | 285 | ("cache", self.cache_backend, lambda: self._ping(self.cache_backend)), |
| 249 | ) | 286 | ) |
| @@ -260,7 +297,7 @@ class SystemContext: | |||
| 260 | async def _stop_resources(self, *, suppress_errors: bool) -> None: | 297 | async def _stop_resources(self, *, suppress_errors: bool) -> None: |
| 261 | errors: list[BaseException] = [] | 298 | errors: list[BaseException] = [] |
| 262 | active = set(self._active_resources) | 299 | active = set(self._active_resources) |
| 263 | - for name in ("cache", "lock", "redis", "db"): | 300 | + for name in ("cache", "lock", "kubernetes", "redis", "db"): |
| 264 | if name not in active: | 301 | if name not in active: |
| 265 | continue | 302 | continue |
| 266 | try: | 303 | try: |
| @@ -286,6 +323,13 @@ class SystemContext: | |||
| 286 | if callable(close): | 323 | if callable(close): |
| 287 | await self._call(close) | 324 | await self._call(close) |
| 288 | return | 325 | return |
| 326 | + if ( | ||
| 327 | + name == "kubernetes" | ||
| 328 | + and self.kubernetes is not None | ||
| 329 | + and self._owns_kubernetes | ||
| 330 | + ): | ||
| 331 | + await self._call(self.kubernetes.close) | ||
| 332 | + return | ||
| 289 | if name == "redis" and self.redis is not None and self._owns_redis: | 333 | if name == "redis" and self.redis is not None and self._owns_redis: |
| 290 | close = getattr(self.redis, "aclose", None) or getattr( | 334 | close = getattr(self.redis, "aclose", None) or getattr( |
| 291 | self.redis, "close", None | 335 | self.redis, "close", None |
| @@ -416,6 +460,7 @@ class SystemContext: | |||
| 416 | etcd_client: Any = None, | 460 | etcd_client: Any = None, |
| 417 | lock_backend: Any = None, | 461 | lock_backend: Any = None, |
| 418 | cache_backend: Any = None, | 462 | cache_backend: Any = None, |
| 463 | + kubernetes: KubernetesOperations | None = None, | ||
| 419 | table_definitions: Iterable[Any] | None = None, | 464 | table_definitions: Iterable[Any] | None = None, |
| 420 | instance_id: str | None = None, | 465 | instance_id: str | None = None, |
| 421 | key_prefix: str | None = None, | 466 | key_prefix: str | None = None, |
| @@ -443,6 +488,7 @@ class SystemContext: | |||
| 443 | etcd_client=etcd_client, | 488 | etcd_client=etcd_client, |
| 444 | lock_backend=lock_backend, | 489 | lock_backend=lock_backend, |
| 445 | cache_backend=cache_backend, | 490 | cache_backend=cache_backend, |
| 491 | + kubernetes=kubernetes, | ||
| 446 | table_definitions=table_definitions, | 492 | table_definitions=table_definitions, |
| 447 | instance_id=instance_id, | 493 | instance_id=instance_id, |
| 448 | ) | 494 | ) |
| @@ -28,6 +28,8 @@ class ErrorCode: | |||
| 28 | REDIS_UNAVAILABLE = "redis_unavailable" | 28 | REDIS_UNAVAILABLE = "redis_unavailable" |
| 29 | CACHE_UNAVAILABLE = "cache_unavailable" | 29 | CACHE_UNAVAILABLE = "cache_unavailable" |
| 30 | LOCK_BACKEND_UNAVAILABLE = "lock_backend_unavailable" | 30 | LOCK_BACKEND_UNAVAILABLE = "lock_backend_unavailable" |
| 31 | + KUBERNETES_UNAVAILABLE = "kubernetes_unavailable" | ||
| 32 | + FORBIDDEN = "forbidden" | ||
| 31 | 33 | ||
| 32 | 34 | ||
| 33 | class FrameworkError(Exception): | 35 | class FrameworkError(Exception): |
| @@ -125,6 +127,18 @@ class CacheUnavailable(FrameworkError): | |||
| 125 | code = ErrorCode.CACHE_UNAVAILABLE | 127 | code = ErrorCode.CACHE_UNAVAILABLE |
| 126 | 128 | ||
| 127 | 129 | ||
| 130 | +class KubernetesUnavailable(FrameworkError): | ||
| 131 | + """Kubernetes operations are absent, closed, or unreachable.""" | ||
| 132 | + | ||
| 133 | + code = ErrorCode.KUBERNETES_UNAVAILABLE | ||
| 134 | + | ||
| 135 | + | ||
| 136 | +class PermissionDenied(FrameworkError): | ||
| 137 | + """The caller or service identity lacks permission for an operation.""" | ||
| 138 | + | ||
| 139 | + code = ErrorCode.FORBIDDEN | ||
| 140 | + | ||
| 141 | + | ||
| 128 | 142 | ||
| 129 | class _HasCode(Protocol): | 143 | class _HasCode(Protocol): |
| 130 | code: str | 144 | code: str |
| @@ -150,6 +164,8 @@ _HTTP_STATUS = { | |||
| 150 | ErrorCode.REDIS_UNAVAILABLE: 503, | 164 | ErrorCode.REDIS_UNAVAILABLE: 503, |
| 151 | ErrorCode.CACHE_UNAVAILABLE: 503, | 165 | ErrorCode.CACHE_UNAVAILABLE: 503, |
| 152 | ErrorCode.LOCK_BACKEND_UNAVAILABLE: 503, | 166 | ErrorCode.LOCK_BACKEND_UNAVAILABLE: 503, |
| 167 | + ErrorCode.KUBERNETES_UNAVAILABLE: 503, | ||
| 168 | + ErrorCode.FORBIDDEN: 403, | ||
| 153 | ErrorCode.INTERNAL: 500, | 169 | ErrorCode.INTERNAL: 500, |
| 154 | } | 170 | } |
| 155 | 171 | ||
| @@ -25,6 +25,9 @@ dependencies = [ | |||
| 25 | etcd = [ | 25 | etcd = [ |
| 26 | "aetcd==1.0.0a4", | 26 | "aetcd==1.0.0a4", |
| 27 | ] | 27 | ] |
| 28 | +kubernetes = [ | ||
| 29 | + "kubernetes_asyncio==35.0.0", | ||
| 30 | +] | ||
| 28 | 31 | ||
| 29 | [tool.uv] | 32 | [tool.uv] |
| 30 | default-groups = ['dev'] | 33 | default-groups = ['dev'] |
| @@ -108,6 +108,7 @@ async def test_system_builder_starts_sqlite_memory_resources_and_reports_readine | |||
| 108 | assert await system.readiness() == { | 108 | assert await system.readiness() == { |
| 109 | "db": True, | 109 | "db": True, |
| 110 | "redis": None, | 110 | "redis": None, |
| 111 | + "kubernetes": None, | ||
| 111 | "lock": True, | 112 | "lock": True, |
| 112 | "cache": True, | 113 | "cache": True, |
| 113 | "ready": True, | 114 | "ready": True, |
| @@ -183,6 +184,17 @@ async def test_start_failure_closes_owned_resources_in_reverse_order(): | |||
| 183 | async def close(self): | 184 | async def close(self): |
| 184 | events.append("lock.close") | 185 | events.append("lock.close") |
| 185 | 186 | ||
| 187 | + class Kubernetes: | ||
| 188 | + async def start(self): | ||
| 189 | + events.append("kubernetes.start") | ||
| 190 | + | ||
| 191 | + async def ping(self): | ||
| 192 | + events.append("kubernetes.ping") | ||
| 193 | + return True | ||
| 194 | + | ||
| 195 | + async def close(self): | ||
| 196 | + events.append("kubernetes.close") | ||
| 197 | + | ||
| 186 | class Cache: | 198 | class Cache: |
| 187 | async def ping(self): | 199 | async def ping(self): |
| 188 | events.append("cache.ping") | 200 | events.append("cache.ping") |
| @@ -194,11 +206,13 @@ async def test_start_failure_closes_owned_resources_in_reverse_order(): | |||
| 194 | system = SystemContext( | 206 | system = SystemContext( |
| 195 | db=Db(), | 207 | db=Db(), |
| 196 | redis=Redis(), | 208 | redis=Redis(), |
| 209 | + kubernetes=Kubernetes(), | ||
| 197 | lock_backend=Lock(), | 210 | lock_backend=Lock(), |
| 198 | cache_backend=Cache(), | 211 | cache_backend=Cache(), |
| 199 | settings=ServiceConfig(deploy_replicas=2), | 212 | settings=ServiceConfig(deploy_replicas=2), |
| 200 | _owns_db=True, | 213 | _owns_db=True, |
| 201 | _owns_redis=True, | 214 | _owns_redis=True, |
| 215 | + _owns_kubernetes=True, | ||
| 202 | _owns_lock_backend=True, | 216 | _owns_lock_backend=True, |
| 203 | _owns_cache_backend=True, | 217 | _owns_cache_backend=True, |
| 204 | ) | 218 | ) |
| @@ -211,10 +225,13 @@ async def test_start_failure_closes_owned_resources_in_reverse_order(): | |||
| 211 | "db.connect", | 225 | "db.connect", |
| 212 | "db.ping", | 226 | "db.ping", |
| 213 | "redis.ping", | 227 | "redis.ping", |
| 228 | + "kubernetes.start", | ||
| 229 | + "kubernetes.ping", | ||
| 214 | "lock.ping", | 230 | "lock.ping", |
| 215 | "cache.ping", | 231 | "cache.ping", |
| 216 | "cache.close", | 232 | "cache.close", |
| 217 | "lock.close", | 233 | "lock.close", |
| 234 | + "kubernetes.close", | ||
| 218 | "redis.close", | 235 | "redis.close", |
| 219 | "db.close", | 236 | "db.close", |
| 220 | ] | 237 | ] |
| @@ -13,9 +13,11 @@ from openjiuwen_runtime.service.errors import ( | |||
| 13 | FrameworkError, | 13 | FrameworkError, |
| 14 | IdempotentConflict, | 14 | IdempotentConflict, |
| 15 | Interrupted, | 15 | Interrupted, |
| 16 | + KubernetesUnavailable, | ||
| 16 | LockLost, | 17 | LockLost, |
| 17 | LockNotAcquired, | 18 | LockNotAcquired, |
| 18 | NotFoundError, | 19 | NotFoundError, |
| 20 | + PermissionDenied, | ||
| 19 | RedisUnavailable, | 21 | RedisUnavailable, |
| 20 | ValidationError, | 22 | ValidationError, |
| 21 | exception_code, | 23 | exception_code, |
| @@ -43,6 +45,8 @@ def test_subclass_codes(): | |||
| 43 | assert DatabaseUnavailable("x").code == ErrorCode.DATABASE_UNAVAILABLE | 45 | assert DatabaseUnavailable("x").code == ErrorCode.DATABASE_UNAVAILABLE |
| 44 | assert RedisUnavailable("x").code == ErrorCode.REDIS_UNAVAILABLE | 46 | assert RedisUnavailable("x").code == ErrorCode.REDIS_UNAVAILABLE |
| 45 | assert CacheUnavailable("x").code == ErrorCode.CACHE_UNAVAILABLE | 47 | assert CacheUnavailable("x").code == ErrorCode.CACHE_UNAVAILABLE |
| 48 | + assert KubernetesUnavailable("x").code == ErrorCode.KUBERNETES_UNAVAILABLE | ||
| 49 | + assert PermissionDenied("x").code == ErrorCode.FORBIDDEN | ||
| 46 | # 都是 FrameworkError 子类 → 中间件可统一捕获 | 50 | # 都是 FrameworkError 子类 → 中间件可统一捕获 |
| 47 | for exc in ( | 51 | for exc in ( |
| 48 | ValidationError(""), | 52 | ValidationError(""), |
| @@ -82,6 +86,8 @@ def test_http_status_mapping(): | |||
| 82 | assert http_status_for(ErrorCode.DATABASE_UNAVAILABLE) == 503 | 86 | assert http_status_for(ErrorCode.DATABASE_UNAVAILABLE) == 503 |
| 83 | assert http_status_for(ErrorCode.REDIS_UNAVAILABLE) == 503 | 87 | assert http_status_for(ErrorCode.REDIS_UNAVAILABLE) == 503 |
| 84 | assert http_status_for(ErrorCode.CACHE_UNAVAILABLE) == 503 | 88 | assert http_status_for(ErrorCode.CACHE_UNAVAILABLE) == 503 |
| 89 | + assert http_status_for(ErrorCode.KUBERNETES_UNAVAILABLE) == 503 | ||
| 90 | + assert http_status_for(ErrorCode.FORBIDDEN) == 403 | ||
| 85 | assert http_status_for(ErrorCode.INTERNAL) == 500 | 91 | assert http_status_for(ErrorCode.INTERNAL) == 500 |
| 86 | # 未知 code → 500(fail-safe) | 92 | # 未知 code → 500(fail-safe) |
| 87 | assert http_status_for("totally-unknown") == 500 | 93 | assert http_status_for("totally-unknown") == 500 |
| @@ -0,0 +1,267 @@ | |||
| 1 | +# coding: utf-8 | ||
| 2 | +# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved | ||
| 3 | + | ||
| 4 | +"""Kubernetes Pod capability tests.""" | ||
| 5 | + | ||
| 6 | +from __future__ import annotations | ||
| 7 | + | ||
| 8 | +from types import SimpleNamespace | ||
| 9 | +from unittest.mock import AsyncMock, patch | ||
| 10 | + | ||
| 11 | +import pytest | ||
| 12 | + | ||
| 13 | +from openjiuwen_runtime.service import ( | ||
| 14 | + ErrorCode, | ||
| 15 | + FakeKubernetesOperations, | ||
| 16 | + FrameworkError, | ||
| 17 | + KubernetesAsyncioOperations, | ||
| 18 | + KubernetesUnavailable, | ||
| 19 | + NotFoundError, | ||
| 20 | + PermissionDenied, | ||
| 21 | + PodCreateSpec, | ||
| 22 | + PodSummary, | ||
| 23 | +) | ||
| 24 | + | ||
| 25 | + | ||
| 26 | + | ||
| 27 | +async def test_fake_kubernetes_lifecycle_and_pod_crud(): | ||
| 28 | + operations = FakeKubernetesOperations("demo") | ||
| 29 | + | ||
| 30 | + assert await operations.ping() is False | ||
| 31 | + await operations.start() | ||
| 32 | + assert await operations.ping() is True | ||
| 33 | + | ||
| 34 | + created = await operations.create_pod( | ||
| 35 | + PodCreateSpec(name="capability-pod-1", image="demo:1") | ||
| 36 | + ) | ||
| 37 | + assert created == PodSummary( | ||
| 38 | + name="capability-pod-1", | ||
| 39 | + namespace="demo", | ||
| 40 | + phase="Running", | ||
| 41 | + ready=True, | ||
| 42 | + image="demo:1", | ||
| 43 | + ) | ||
| 44 | + assert await operations.get_pod("capability-pod-1") == created | ||
| 45 | + | ||
| 46 | + with pytest.raises(FrameworkError) as conflict: | ||
| 47 | + await operations.create_pod( | ||
| 48 | + PodCreateSpec(name="capability-pod-1", image="demo:2") | ||
| 49 | + ) | ||
| 50 | + assert conflict.value.code == ErrorCode.CONFLICT | ||
| 51 | + | ||
| 52 | + deleted = await operations.delete_pod("capability-pod-1") | ||
| 53 | + absent = await operations.delete_pod("capability-pod-1") | ||
| 54 | + assert deleted.state == "delete_requested" | ||
| 55 | + assert absent.state == "already_absent" | ||
| 56 | + assert await operations.get_pod("capability-pod-1") is None | ||
| 57 | + | ||
| 58 | + await operations.close() | ||
| 59 | + assert await operations.ping() is False | ||
| 60 | + with pytest.raises(KubernetesUnavailable): | ||
| 61 | + await operations.get_pod("capability-pod-1") | ||
| 62 | + | ||
| 63 | + | ||
| 64 | +def _pod( | ||
| 65 | + *, | ||
| 66 | + labels: dict[str, str] | None = None, | ||
| 67 | + deletion_timestamp=None, | ||
| 68 | + uid: str = "uid-1", | ||
| 69 | +): | ||
| 70 | + return SimpleNamespace( | ||
| 71 | + metadata=SimpleNamespace( | ||
| 72 | + name="capability-pod-1", | ||
| 73 | + namespace="demo", | ||
| 74 | + labels=labels | ||
| 75 | + or {"app.kubernetes.io/managed-by": "openjiuwen-service-demo"}, | ||
| 76 | + deletion_timestamp=deletion_timestamp, | ||
| 77 | + uid=uid, | ||
| 78 | + ), | ||
| 79 | + spec=SimpleNamespace( | ||
| 80 | + containers=[SimpleNamespace(image="demo:1")], | ||
| 81 | + ), | ||
| 82 | + status=SimpleNamespace( | ||
| 83 | + phase="Running", | ||
| 84 | + conditions=[SimpleNamespace(type="Ready", status="True")], | ||
| 85 | + container_statuses=[SimpleNamespace(ready=True)], | ||
| 86 | + ), | ||
| 87 | + ) | ||
| 88 | + | ||
| 89 | + | ||
| 90 | + | ||
| 91 | +def test_real_adapter_converts_pod_status_and_terminating_phase(): | ||
| 92 | + operations = KubernetesAsyncioOperations("demo") | ||
| 93 | + | ||
| 94 | + running = operations._to_summary(_pod()) | ||
| 95 | + terminating = operations._to_summary(_pod(deletion_timestamp="now")) | ||
| 96 | + | ||
| 97 | + assert running == PodSummary( | ||
| 98 | + name="capability-pod-1", | ||
| 99 | + namespace="demo", | ||
| 100 | + phase="Running", | ||
| 101 | + ready=True, | ||
| 102 | + image="demo:1", | ||
| 103 | + ) | ||
| 104 | + assert terminating.phase == "Terminating" | ||
| 105 | + | ||
| 106 | + | ||
| 107 | + | ||
| 108 | + | ||
| 109 | + ("status", "error_type", "code"), | ||
| 110 | + [ | ||
| 111 | + (401, PermissionDenied, ErrorCode.FORBIDDEN), | ||
| 112 | + (403, PermissionDenied, ErrorCode.FORBIDDEN), | ||
| 113 | + (404, NotFoundError, ErrorCode.NOT_FOUND), | ||
| 114 | + (409, FrameworkError, ErrorCode.CONFLICT), | ||
| 115 | + (429, KubernetesUnavailable, ErrorCode.KUBERNETES_UNAVAILABLE), | ||
| 116 | + (503, KubernetesUnavailable, ErrorCode.KUBERNETES_UNAVAILABLE), | ||
| 117 | + (None, KubernetesUnavailable, ErrorCode.KUBERNETES_UNAVAILABLE), | ||
| 118 | + ], | ||
| 119 | +) | ||
| 120 | +def test_real_adapter_maps_api_errors(status, error_type, code): | ||
| 121 | + exc = RuntimeError("api unavailable") | ||
| 122 | + exc.status = status | ||
| 123 | + exc.reason = "test reason" | ||
| 124 | + | ||
| 125 | + mapped = KubernetesAsyncioOperations._map_api_exception( | ||
| 126 | + exc, | ||
| 127 | + operation="test", | ||
| 128 | + ) | ||
| 129 | + | ||
| 130 | + assert isinstance(mapped, error_type) | ||
| 131 | + assert mapped.code == code | ||
| 132 | + assert "test reason" in mapped.message | ||
| 133 | + | ||
| 134 | + | ||
| 135 | + | ||
| 136 | +async def test_real_adapter_checks_labels_and_delete_state(): | ||
| 137 | + operations = KubernetesAsyncioOperations("demo") | ||
| 138 | + operations._core_api = SimpleNamespace( | ||
| 139 | + read_namespaced_pod=AsyncMock( | ||
| 140 | + side_effect=[ | ||
| 141 | + _pod(labels={"owner": "other"}), | ||
| 142 | + _pod(labels={"owner": "other"}), | ||
| 143 | + _pod(deletion_timestamp="now"), | ||
| 144 | + ] | ||
| 145 | + ) | ||
| 146 | + ) | ||
| 147 | + | ||
| 148 | + assert await operations.get_pod("capability-pod-1") is None | ||
| 149 | + with pytest.raises(NotFoundError): | ||
| 150 | + await operations.delete_pod("capability-pod-1") | ||
| 151 | + result = await operations.delete_pod("capability-pod-1") | ||
| 152 | + assert result.state == "deletion_in_progress" | ||
| 153 | + | ||
| 154 | + | ||
| 155 | + | ||
| 156 | +async def test_real_adapter_delete_uses_uid_precondition(): | ||
| 157 | + client = pytest.importorskip("kubernetes_asyncio.client") | ||
| 158 | + operations = KubernetesAsyncioOperations("demo") | ||
| 159 | + operations._client_module = client | ||
| 160 | + operations._core_api = SimpleNamespace( | ||
| 161 | + read_namespaced_pod=AsyncMock(return_value=_pod(uid="stable-uid")), | ||
| 162 | + delete_namespaced_pod=AsyncMock(), | ||
| 163 | + ) | ||
| 164 | + | ||
| 165 | + result = await operations.delete_pod("capability-pod-1") | ||
| 166 | + | ||
| 167 | + assert result.state == "delete_requested" | ||
| 168 | + call = operations._core_api.delete_namespaced_pod.await_args.kwargs | ||
| 169 | + assert call["body"].preconditions.uid == "stable-uid" | ||
| 170 | + assert call["grace_period_seconds"] == 0 | ||
| 171 | + assert call["_request_timeout"] == 10 | ||
| 172 | + | ||
| 173 | + | ||
| 174 | + | ||
| 175 | +async def test_real_adapter_create_uses_restricted_pod_template(): | ||
| 176 | + client = pytest.importorskip("kubernetes_asyncio.client") | ||
| 177 | + operations = KubernetesAsyncioOperations("demo") | ||
| 178 | + operations._client_module = client | ||
| 179 | + operations._core_api = SimpleNamespace( | ||
| 180 | + create_namespaced_pod=AsyncMock(return_value=_pod()), | ||
| 181 | + ) | ||
| 182 | + | ||
| 183 | + await operations.create_pod( | ||
| 184 | + PodCreateSpec(name="capability-pod-1", image="demo:1") | ||
| 185 | + ) | ||
| 186 | + | ||
| 187 | + call = operations._core_api.create_namespaced_pod.await_args.kwargs | ||
| 188 | + body = call["body"] | ||
| 189 | + container = body.spec.containers[0] | ||
| 190 | + security = container.security_context | ||
| 191 | + assert body.metadata.labels == operations.DEFAULT_LABELS | ||
| 192 | + assert body.spec.restart_policy == "Never" | ||
| 193 | + assert body.spec.automount_service_account_token is False | ||
| 194 | + assert container.name == "demo" | ||
| 195 | + assert container.image == "demo:1" | ||
| 196 | + assert container.image_pull_policy == "IfNotPresent" | ||
| 197 | + assert security.allow_privilege_escalation is False | ||
| 198 | + assert security.run_as_non_root is True | ||
| 199 | + assert security.capabilities.drop == ["ALL"] | ||
| 200 | + assert security.seccomp_profile.type == "RuntimeDefault" | ||
| 201 | + assert container.resources.requests == {"cpu": "10m", "memory": "16Mi"} | ||
| 202 | + assert container.resources.limits == {"cpu": "100m", "memory": "64Mi"} | ||
| 203 | + | ||
| 204 | + | ||
| 205 | + | ||
| 206 | +async def test_real_adapter_start_falls_back_to_kubeconfig_once(): | ||
| 207 | + pytest.importorskip("kubernetes_asyncio") | ||
| 208 | + from kubernetes_asyncio.config.config_exception import ConfigException | ||
| 209 | + | ||
| 210 | + api_client = SimpleNamespace(close=AsyncMock()) | ||
| 211 | + core_api = object() | ||
| 212 | + load_kube_config = AsyncMock() | ||
| 213 | + with ( | ||
| 214 | + patch( | ||
| 215 | + "kubernetes_asyncio.config.load_incluster_config", | ||
| 216 | + side_effect=ConfigException("not in cluster"), | ||
| 217 | + ) as load_incluster, | ||
| 218 | + patch( | ||
| 219 | + "kubernetes_asyncio.config.load_kube_config", | ||
| 220 | + load_kube_config, | ||
| 221 | + ), | ||
| 222 | + patch("kubernetes_asyncio.client.ApiClient", return_value=api_client), | ||
| 223 | + patch("kubernetes_asyncio.client.CoreV1Api", return_value=core_api), | ||
| 224 | + ): | ||
| 225 | + operations = KubernetesAsyncioOperations( | ||
| 226 | + "demo", | ||
| 227 | + kubeconfig="C:/kube/config", | ||
| 228 | + ) | ||
| 229 | + await operations.start() | ||
| 230 | + await operations.start() | ||
| 231 | + | ||
| 232 | + load_incluster.assert_called_once_with() | ||
| 233 | + load_kube_config.assert_awaited_once_with(config_file="C:/kube/config") | ||
| 234 | + assert operations._api_client is api_client | ||
| 235 | + assert operations._core_api is core_api | ||
| 236 | + | ||
| 237 | + | ||
| 238 | + | ||
| 239 | +async def test_real_adapter_returns_absent_for_api_404(): | ||
| 240 | + exc = RuntimeError("not found") | ||
| 241 | + exc.status = 404 | ||
| 242 | + operations = KubernetesAsyncioOperations("demo") | ||
| 243 | + operations._core_api = SimpleNamespace( | ||
| 244 | + read_namespaced_pod=AsyncMock(side_effect=exc), | ||
| 245 | + ) | ||
| 246 | + | ||
| 247 | + assert await operations.get_pod("missing") is None | ||
| 248 | + result = await operations.delete_pod("missing") | ||
| 249 | + assert result.state == "already_absent" | ||
| 250 | + | ||
| 251 | + | ||
| 252 | + | ||
| 253 | +async def test_real_adapter_ping_reuses_core_api_and_close_is_idempotent(): | ||
| 254 | + api_client = SimpleNamespace(close=AsyncMock()) | ||
| 255 | + operations = KubernetesAsyncioOperations("demo") | ||
| 256 | + operations._api_client = api_client | ||
| 257 | + operations._core_api = SimpleNamespace(list_namespaced_pod=AsyncMock()) | ||
| 258 | + | ||
| 259 | + assert await operations.ping() is True | ||
| 260 | + operations._core_api.list_namespaced_pod.assert_awaited_once_with( | ||
| 261 | + namespace="demo", | ||
| 262 | + limit=1, | ||
| 263 | + _request_timeout=10, | ||
| 264 | + ) | ||
| 265 | + await operations.close() | ||
| 266 | + await operations.close() | ||
| 267 | + api_client.close.assert_awaited_once() | ||
| @@ -0,0 +1,518 @@ | |||
| 1 | +# coding: utf-8 | ||
| 2 | +# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved | ||
| 3 | + | ||
| 4 | +"""Tests for the deployable simple capabilities demo.""" | ||
| 5 | + | ||
| 6 | +from __future__ import annotations | ||
| 7 | + | ||
| 8 | +import importlib | ||
| 9 | +import os | ||
| 10 | +from dataclasses import dataclass | ||
| 11 | +from pathlib import Path | ||
| 12 | +from unittest.mock import patch | ||
| 13 | + | ||
| 14 | +import pytest | ||
| 15 | +from starlette.testclient import TestClient | ||
| 16 | + | ||
| 17 | + | ||
| 18 | +_IMPORT_ENV = { | ||
| 19 | + "DEMO_ENVIRONMENT": "test", | ||
| 20 | + "DEMO_SERVICE_LABEL": "simple-capabilities-test", | ||
| 21 | + "DEMO_REDIS_MODE": "fake", | ||
| 22 | + "DEMO_REDIS_DEFAULT_TTL_SECONDS": "25", | ||
| 23 | + "DEMO_KUBERNETES_MODE": "fake", | ||
| 24 | + "DEMO_KUBERNETES_NAMESPACE": "simple-capabilities-demo", | ||
| 25 | + "DEMO_KUBERNETES_POD_IMAGE": "demo-pod:test", | ||
| 26 | + "OPENJIUWEN_SERVICE_REDIS_URL": "disabled", | ||
| 27 | +} | ||
| 28 | + | ||
| 29 | + | ||
| 30 | + | ||
| 31 | +class _InvalidRedisConfigCase: | ||
| 32 | + mode: str | ||
| 33 | + ttl: str | ||
| 34 | + redis_url: str | ||
| 35 | + message: str | ||
| 36 | + | ||
| 37 | + | ||
| 38 | + | ||
| 39 | +def demo_module(): | ||
| 40 | + with patch.dict(os.environ, _IMPORT_ENV, clear=False): | ||
| 41 | + yield importlib.import_module("examples.simple_capabilities_app") | ||
| 42 | + | ||
| 43 | + | ||
| 44 | + | ||
| 45 | +def demo_client(demo_module, monkeypatch, tmp_path: Path): | ||
| 46 | + database_path = tmp_path / "simple-capabilities.db" | ||
| 47 | + environment = { | ||
| 48 | + "OPENJIUWEN_SERVICE_DB_TYPE": "sqlite", | ||
| 49 | + "OPENJIUWEN_SERVICE_DB_NAME": str(database_path), | ||
| 50 | + "OPENJIUWEN_SERVICE_REDIS_URL": "disabled", | ||
| 51 | + "OPENJIUWEN_SERVICE_REDIS_KEY_PREFIX": "simple-capabilities-test", | ||
| 52 | + "OPENJIUWEN_SERVICE_LOCK_BACKEND": "memory", | ||
| 53 | + "OPENJIUWEN_SERVICE_CACHE_BACKEND": "memory", | ||
| 54 | + "OPENJIUWEN_SERVICE_DEPLOY_REPLICAS": "1", | ||
| 55 | + "OPENJIUWEN_SERVICE_REQUEST_TIMEOUT_SECONDS": "30", | ||
| 56 | + } | ||
| 57 | + for name, value in environment.items(): | ||
| 58 | + monkeypatch.setenv(name, value) | ||
| 59 | + | ||
| 60 | + config = demo_module.DemoConfig( | ||
| 61 | + environment="test", | ||
| 62 | + service_label="simple-capabilities-test", | ||
| 63 | + redis_mode="fake", | ||
| 64 | + redis_default_ttl_seconds=25, | ||
| 65 | + kubernetes_mode="fake", | ||
| 66 | + kubernetes_namespace="simple-capabilities-demo", | ||
| 67 | + kubernetes_pod_image="demo-pod:test", | ||
| 68 | + ) | ||
| 69 | + application = demo_module.create_app(config) | ||
| 70 | + with TestClient(application.asgi) as client: | ||
| 71 | + yield client | ||
| 72 | + | ||
| 73 | + | ||
| 74 | +def _envelope( | ||
| 75 | + msg_type: str, | ||
| 76 | + rawdata: dict, | ||
| 77 | + *, | ||
| 78 | + request_id: str, | ||
| 79 | + full_metadata: bool = False, | ||
| 80 | +) -> dict: | ||
| 81 | + metadata = {"request_id": request_id} | ||
| 82 | + if full_metadata: | ||
| 83 | + metadata.update( | ||
| 84 | + { | ||
| 85 | + "user_id": "demo-user", | ||
| 86 | + "chat_id": "demo-chat", | ||
| 87 | + "session_id": "demo-session", | ||
| 88 | + "bot_id": "demo-bot", | ||
| 89 | + "channel": "swagger", | ||
| 90 | + "timestamp": 1786500000, | ||
| 91 | + "trace_id": "trace-1", | ||
| 92 | + "instance_id": "demo-instance", | ||
| 93 | + "extra": {"tenant": "demo"}, | ||
| 94 | + } | ||
| 95 | + ) | ||
| 96 | + return { | ||
| 97 | + "type": msg_type, | ||
| 98 | + "metadata": metadata, | ||
| 99 | + "rawdata": rawdata, | ||
| 100 | + "version": "1", | ||
| 101 | + } | ||
| 102 | + | ||
| 103 | + | ||
| 104 | + | ||
| 105 | + | ||
| 106 | + ("mode", "redis_url"), | ||
| 107 | + [("fake", "disabled"), ("real", "redis://redis.internal:6379/2")], | ||
| 108 | +) | ||
| 109 | +def test_demo_config_reads_supported_modes( | ||
| 110 | + demo_module, | ||
| 111 | + monkeypatch, | ||
| 112 | + mode: str, | ||
| 113 | + redis_url: str, | ||
| 114 | +): | ||
| 115 | + monkeypatch.setenv("DEMO_ENVIRONMENT", "test") | ||
| 116 | + monkeypatch.setenv("DEMO_SERVICE_LABEL", "demo-test") | ||
| 117 | + monkeypatch.setenv("DEMO_REDIS_MODE", mode) | ||
| 118 | + monkeypatch.setenv("DEMO_REDIS_DEFAULT_TTL_SECONDS", "45") | ||
| 119 | + monkeypatch.setenv("OPENJIUWEN_SERVICE_REDIS_URL", redis_url) | ||
| 120 | + monkeypatch.setenv("DEMO_KUBERNETES_MODE", mode) | ||
| 121 | + monkeypatch.setenv("DEMO_KUBERNETES_NAMESPACE", "simple-capabilities-demo") | ||
| 122 | + monkeypatch.setenv("DEMO_KUBERNETES_POD_IMAGE", "demo-pod:test") | ||
| 123 | + | ||
| 124 | + config = demo_module.DemoConfig.from_env() | ||
| 125 | + | ||
| 126 | + assert config.redis_mode == mode | ||
| 127 | + assert config.redis_default_ttl_seconds == 45 | ||
| 128 | + assert config.kubernetes_mode == mode | ||
| 129 | + | ||
| 130 | + | ||
| 131 | + | ||
| 132 | + | ||
| 133 | + "case", | ||
| 134 | + [ | ||
| 135 | + _InvalidRedisConfigCase("other", "30", "disabled", "DEMO_REDIS_MODE"), | ||
| 136 | + _InvalidRedisConfigCase("fake", "0", "disabled", "positive integer"), | ||
| 137 | + _InvalidRedisConfigCase( | ||
| 138 | + "fake", "invalid", "disabled", "positive integer" | ||
| 139 | + ), | ||
| 140 | + _InvalidRedisConfigCase( | ||
| 141 | + "fake", | ||
| 142 | + "30", | ||
| 143 | + "redis://localhost:6379/0", | ||
| 144 | + "fake Redis mode", | ||
| 145 | + ), | ||
| 146 | + _InvalidRedisConfigCase("real", "30", "disabled", "real Redis mode"), | ||
| 147 | + _InvalidRedisConfigCase("real", "30", "not-a-url", "valid Redis URL"), | ||
| 148 | + ], | ||
| 149 | +) | ||
| 150 | +def test_demo_config_rejects_invalid_values( | ||
| 151 | + demo_module, | ||
| 152 | + monkeypatch, | ||
| 153 | + case: _InvalidRedisConfigCase, | ||
| 154 | +): | ||
| 155 | + monkeypatch.setenv("DEMO_ENVIRONMENT", "test") | ||
| 156 | + monkeypatch.setenv("DEMO_SERVICE_LABEL", "demo-test") | ||
| 157 | + monkeypatch.setenv("DEMO_REDIS_MODE", case.mode) | ||
| 158 | + monkeypatch.setenv("DEMO_REDIS_DEFAULT_TTL_SECONDS", case.ttl) | ||
| 159 | + monkeypatch.setenv("OPENJIUWEN_SERVICE_REDIS_URL", case.redis_url) | ||
| 160 | + monkeypatch.setenv("DEMO_KUBERNETES_MODE", "fake") | ||
| 161 | + monkeypatch.setenv("DEMO_KUBERNETES_NAMESPACE", "simple-capabilities-demo") | ||
| 162 | + monkeypatch.setenv("DEMO_KUBERNETES_POD_IMAGE", "demo-pod:test") | ||
| 163 | + | ||
| 164 | + with pytest.raises(ValueError, match=case.message): | ||
| 165 | + demo_module.DemoConfig.from_env() | ||
| 166 | + | ||
| 167 | + | ||
| 168 | + | ||
| 169 | + | ||
| 170 | + ("mode", "namespace", "image", "message"), | ||
| 171 | + [ | ||
| 172 | + ("other", "simple-capabilities-demo", "demo:1", "DEMO_KUBERNETES_MODE"), | ||
| 173 | + ("fake", "", "demo:1", "DEMO_KUBERNETES_NAMESPACE"), | ||
| 174 | + ("fake", "UPPERCASE", "demo:1", "DEMO_KUBERNETES_NAMESPACE"), | ||
| 175 | + ("fake", "simple-capabilities-demo", "", "DEMO_KUBERNETES_POD_IMAGE"), | ||
| 176 | + ], | ||
| 177 | +) | ||
| 178 | +def test_demo_config_rejects_invalid_kubernetes_values( | ||
| 179 | + demo_module, | ||
| 180 | + mode: str, | ||
| 181 | + namespace: str, | ||
| 182 | + image: str, | ||
| 183 | + message: str, | ||
| 184 | +): | ||
| 185 | + with pytest.raises(ValueError, match=message): | ||
| 186 | + demo_module.DemoConfig( | ||
| 187 | + environment="test", | ||
| 188 | + service_label="demo-test", | ||
| 189 | + redis_mode="fake", | ||
| 190 | + redis_default_ttl_seconds=30, | ||
| 191 | + kubernetes_mode=mode, | ||
| 192 | + kubernetes_namespace=namespace, | ||
| 193 | + kubernetes_pod_image=image, | ||
| 194 | + ) | ||
| 195 | + | ||
| 196 | + | ||
| 197 | + | ||
| 198 | +def test_database_create_update_read_and_missing(demo_client: TestClient): | ||
| 199 | + created = demo_client.post( | ||
| 200 | + "/api/db/write", | ||
| 201 | + json=_envelope( | ||
| 202 | + "db/write", | ||
| 203 | + {"id": "record-1", "value": "first value"}, | ||
| 204 | + request_id="db-write-1", | ||
| 205 | + ), | ||
| 206 | + ) | ||
| 207 | + updated = demo_client.post( | ||
| 208 | + "/api/db/write", | ||
| 209 | + json=_envelope( | ||
| 210 | + "db/write", | ||
| 211 | + {"id": "record-1", "value": "final value"}, | ||
| 212 | + request_id="db-write-2", | ||
| 213 | + ), | ||
| 214 | + ) | ||
| 215 | + read = demo_client.post( | ||
| 216 | + "/api/db/read", | ||
| 217 | + json=_envelope( | ||
| 218 | + "db/read", | ||
| 219 | + {"id": "record-1"}, | ||
| 220 | + request_id="db-read-1", | ||
| 221 | + ), | ||
| 222 | + ) | ||
| 223 | + missing = demo_client.post( | ||
| 224 | + "/api/db/read", | ||
| 225 | + json=_envelope( | ||
| 226 | + "db/read", | ||
| 227 | + {"id": "missing"}, | ||
| 228 | + request_id="db-read-2", | ||
| 229 | + ), | ||
| 230 | + ) | ||
| 231 | + | ||
| 232 | + assert created.status_code == 200 | ||
| 233 | + assert created.json()["rawdata"]["operation"] == "created" | ||
| 234 | + assert created.json()["rawdata"]["record"]["updated_at"].endswith("Z") | ||
| 235 | + assert updated.status_code == 200 | ||
| 236 | + assert updated.json()["rawdata"]["operation"] == "updated" | ||
| 237 | + assert read.status_code == 200 | ||
| 238 | + assert read.json()["rawdata"]["record"]["value"] == "final value" | ||
| 239 | + assert missing.status_code == 404 | ||
| 240 | + assert missing.json()["error_code"] == "not_found" | ||
| 241 | + | ||
| 242 | + | ||
| 243 | + | ||
| 244 | +def test_redis_write_read_ttl_defaults_and_missing(demo_client: TestClient): | ||
| 245 | + explicit = demo_client.post( | ||
| 246 | + "/api/redis/write", | ||
| 247 | + json=_envelope( | ||
| 248 | + "redis/write", | ||
| 249 | + {"key": "explicit", "value": "hello", "ttl_seconds": 120}, | ||
| 250 | + request_id="redis-write-1", | ||
| 251 | + ), | ||
| 252 | + ) | ||
| 253 | + defaulted = demo_client.post( | ||
| 254 | + "/api/redis/write", | ||
| 255 | + json=_envelope( | ||
| 256 | + "redis/write", | ||
| 257 | + {"key": "default", "value": "world"}, | ||
| 258 | + request_id="redis-write-2", | ||
| 259 | + ), | ||
| 260 | + ) | ||
| 261 | + read = demo_client.post( | ||
| 262 | + "/api/redis/read", | ||
| 263 | + json=_envelope( | ||
| 264 | + "redis/read", | ||
| 265 | + {"key": "explicit"}, | ||
| 266 | + request_id="redis-read-1", | ||
| 267 | + ), | ||
| 268 | + ) | ||
| 269 | + missing = demo_client.post( | ||
| 270 | + "/api/redis/read", | ||
| 271 | + json=_envelope( | ||
| 272 | + "redis/read", | ||
| 273 | + {"key": "missing"}, | ||
| 274 | + request_id="redis-read-2", | ||
| 275 | + ), | ||
| 276 | + ) | ||
| 277 | + | ||
| 278 | + redis = demo_client.app.state.sysctx.redis | ||
| 279 | + explicit_ttl = demo_client.portal.call( | ||
| 280 | + redis.ttl, "simple-capabilities-test:kv:explicit" | ||
| 281 | + ) | ||
| 282 | + default_ttl = demo_client.portal.call( | ||
| 283 | + redis.ttl, "simple-capabilities-test:kv:default" | ||
| 284 | + ) | ||
| 285 | + | ||
| 286 | + assert explicit.status_code == 200 | ||
| 287 | + assert explicit.json()["rawdata"]["ttl_seconds"] == 120 | ||
| 288 | + assert 0 < explicit_ttl <= 120 | ||
| 289 | + assert defaulted.json()["rawdata"]["ttl_seconds"] == 25 | ||
| 290 | + assert 0 < default_ttl <= 25 | ||
| 291 | + assert read.json()["rawdata"] == { | ||
| 292 | + "key": "explicit", | ||
| 293 | + "found": True, | ||
| 294 | + "value": "hello", | ||
| 295 | + "backend_mode": "fake", | ||
| 296 | + } | ||
| 297 | + assert missing.status_code == 200 | ||
| 298 | + assert missing.json()["rawdata"]["found"] is False | ||
| 299 | + assert missing.json()["rawdata"]["value"] is None | ||
| 300 | + | ||
| 301 | + | ||
| 302 | + | ||
| 303 | +def test_envelope_inspect_returns_envelope_context_and_service( | ||
| 304 | + demo_client: TestClient, | ||
| 305 | +): | ||
| 306 | + response = demo_client.post( | ||
| 307 | + "/api/envelope/inspect", | ||
| 308 | + json=_envelope( | ||
| 309 | + "envelope/inspect", | ||
| 310 | + { | ||
| 311 | + "message": "inspect this envelope", | ||
| 312 | + "attributes": {"source": "swagger"}, | ||
| 313 | + }, | ||
| 314 | + request_id="inspect-1", | ||
| 315 | + full_metadata=True, | ||
| 316 | + ), | ||
| 317 | + ) | ||
| 318 | + | ||
| 319 | + assert response.status_code == 200 | ||
| 320 | + result = response.json()["rawdata"] | ||
| 321 | + assert result["envelope"] == { | ||
| 322 | + "type": "envelope/inspect", | ||
| 323 | + "version": "1", | ||
| 324 | + "metadata": { | ||
| 325 | + "request_id": "inspect-1", | ||
| 326 | + "user_id": "demo-user", | ||
| 327 | + "chat_id": "demo-chat", | ||
| 328 | + "session_id": "demo-session", | ||
| 329 | + "bot_id": "demo-bot", | ||
| 330 | + "channel": "swagger", | ||
| 331 | + "timestamp": 1786500000.0, | ||
| 332 | + "trace_id": "trace-1", | ||
| 333 | + "instance_id": "demo-instance", | ||
| 334 | + "extra": {"tenant": "demo"}, | ||
| 335 | + }, | ||
| 336 | + "rawdata": { | ||
| 337 | + "message": "inspect this envelope", | ||
| 338 | + "attributes": {"source": "swagger"}, | ||
| 339 | + }, | ||
| 340 | + } | ||
| 341 | + assert result["context"] == { | ||
| 342 | + "msg_type": "envelope/inspect", | ||
| 343 | + "request_id": "inspect-1", | ||
| 344 | + "user_id": "demo-user", | ||
| 345 | + "chat_id": "demo-chat", | ||
| 346 | + "session_id": "demo-session", | ||
| 347 | + "bot_id": "demo-bot", | ||
| 348 | + "channel": "swagger", | ||
| 349 | + "trace_id": "trace-1", | ||
| 350 | + "instance_id": "demo-instance", | ||
| 351 | + "replica_id": result["context"]["replica_id"], | ||
| 352 | + } | ||
| 353 | + assert result["context"]["replica_id"] | ||
| 354 | + assert result["service"] == { | ||
| 355 | + "environment": "test", | ||
| 356 | + "service_label": "simple-capabilities-test", | ||
| 357 | + "redis_mode": "fake", | ||
| 358 | + } | ||
| 359 | + | ||
| 360 | + | ||
| 361 | + | ||
| 362 | +def test_kubernetes_pod_create_read_conflict_delete_and_missing( | ||
| 363 | + demo_client: TestClient, | ||
| 364 | +): | ||
| 365 | + created = demo_client.post( | ||
| 366 | + "/api/k8s/pod/create", | ||
| 367 | + json=_envelope( | ||
| 368 | + "k8s/pod/create", | ||
| 369 | + {"name": "capability-pod-1"}, | ||
| 370 | + request_id="pod-create-1", | ||
| 371 | + ), | ||
| 372 | + ) | ||
| 373 | + read = demo_client.post( | ||
| 374 | + "/api/k8s/pod/read", | ||
| 375 | + json=_envelope( | ||
| 376 | + "k8s/pod/read", | ||
| 377 | + {"name": "capability-pod-1"}, | ||
| 378 | + request_id="pod-read-1", | ||
| 379 | + ), | ||
| 380 | + ) | ||
| 381 | + conflict = demo_client.post( | ||
| 382 | + "/api/k8s/pod/create", | ||
| 383 | + json=_envelope( | ||
| 384 | + "k8s/pod/create", | ||
| 385 | + {"name": "capability-pod-1"}, | ||
| 386 | + request_id="pod-create-2", | ||
| 387 | + ), | ||
| 388 | + ) | ||
| 389 | + deleted = demo_client.post( | ||
| 390 | + "/api/k8s/pod/delete", | ||
| 391 | + json=_envelope( | ||
| 392 | + "k8s/pod/delete", | ||
| 393 | + {"name": "capability-pod-1"}, | ||
| 394 | + request_id="pod-delete-1", | ||
| 395 | + ), | ||
| 396 | + ) | ||
| 397 | + absent = demo_client.post( | ||
| 398 | + "/api/k8s/pod/delete", | ||
| 399 | + json=_envelope( | ||
| 400 | + "k8s/pod/delete", | ||
| 401 | + {"name": "capability-pod-1"}, | ||
| 402 | + request_id="pod-delete-2", | ||
| 403 | + ), | ||
| 404 | + ) | ||
| 405 | + missing = demo_client.post( | ||
| 406 | + "/api/k8s/pod/read", | ||
| 407 | + json=_envelope( | ||
| 408 | + "k8s/pod/read", | ||
| 409 | + {"name": "capability-pod-1"}, | ||
| 410 | + request_id="pod-read-2", | ||
| 411 | + ), | ||
| 412 | + ) | ||
| 413 | + | ||
| 414 | + assert created.status_code == 200 | ||
| 415 | + assert created.json()["rawdata"] == { | ||
| 416 | + "operation": "created", | ||
| 417 | + "pod": { | ||
| 418 | + "name": "capability-pod-1", | ||
| 419 | + "namespace": "simple-capabilities-demo", | ||
| 420 | + "phase": "Running", | ||
| 421 | + "ready": True, | ||
| 422 | + "image": "demo-pod:test", | ||
| 423 | + }, | ||
| 424 | + } | ||
| 425 | + assert read.status_code == 200 | ||
| 426 | + assert read.json()["rawdata"]["pod"] == created.json()["rawdata"]["pod"] | ||
| 427 | + assert conflict.status_code == 409 | ||
| 428 | + assert conflict.json()["error_code"] == "conflict" | ||
| 429 | + assert deleted.json()["rawdata"]["state"] == "delete_requested" | ||
| 430 | + assert absent.json()["rawdata"]["state"] == "already_absent" | ||
| 431 | + assert missing.status_code == 404 | ||
| 432 | + assert missing.json()["error_code"] == "not_found" | ||
| 433 | + | ||
| 434 | + | ||
| 435 | + | ||
| 436 | +def test_kubernetes_pod_name_validation_and_envelope_type(demo_client: TestClient): | ||
| 437 | + invalid = demo_client.post( | ||
| 438 | + "/api/k8s/pod/create", | ||
| 439 | + json=_envelope( | ||
| 440 | + "k8s/pod/create", | ||
| 441 | + {"name": "Invalid_Pod"}, | ||
| 442 | + request_id="pod-invalid-1", | ||
| 443 | + ), | ||
| 444 | + ) | ||
| 445 | + wrong_type = demo_client.post( | ||
| 446 | + "/api/k8s/pod/create", | ||
| 447 | + json=_envelope( | ||
| 448 | + "k8s/pod/read", | ||
| 449 | + {"name": "capability-pod-2"}, | ||
| 450 | + request_id="pod-invalid-2", | ||
| 451 | + ), | ||
| 452 | + ) | ||
| 453 | + | ||
| 454 | + assert invalid.status_code == 400 | ||
| 455 | + assert invalid.json()["error_code"] == "validation" | ||
| 456 | + assert wrong_type.status_code == 422 | ||
| 457 | + | ||
| 458 | + long_label = demo_client.post( | ||
| 459 | + "/api/k8s/pod/create", | ||
| 460 | + json=_envelope( | ||
| 461 | + "k8s/pod/create", | ||
| 462 | + {"name": f"{'a' * 64}.valid"}, | ||
| 463 | + request_id="pod-invalid-3", | ||
| 464 | + ), | ||
| 465 | + ) | ||
| 466 | + assert long_label.status_code == 400 | ||
| 467 | + assert long_label.json()["error_code"] == "validation" | ||
| 468 | + | ||
| 469 | + | ||
| 470 | + | ||
| 471 | +def test_docs_openapi_and_validation_error(demo_client: TestClient): | ||
| 472 | + docs = demo_client.get("/docs") | ||
| 473 | + openapi_response = demo_client.get("/openapi.json") | ||
| 474 | + schema = openapi_response.json() | ||
| 475 | + expected = { | ||
| 476 | + "/api/db/write": "db/write", | ||
| 477 | + "/api/db/read": "db/read", | ||
| 478 | + "/api/redis/write": "redis/write", | ||
| 479 | + "/api/redis/read": "redis/read", | ||
| 480 | + "/api/envelope/inspect": "envelope/inspect", | ||
| 481 | + "/api/k8s/pod/read": "k8s/pod/read", | ||
| 482 | + "/api/k8s/pod/create": "k8s/pod/create", | ||
| 483 | + "/api/k8s/pod/delete": "k8s/pod/delete", | ||
| 484 | + } | ||
| 485 | + | ||
| 486 | + assert docs.status_code == 200 | ||
| 487 | + assert openapi_response.status_code == 200 | ||
| 488 | + for path, msg_type in expected.items(): | ||
| 489 | + operation = schema["paths"][path] | ||
| 490 | + assert set(operation) == {"post"} | ||
| 491 | + request_ref = operation["post"]["requestBody"]["content"]["application/json"][ | ||
| 492 | + "schema" | ||
| 493 | + ]["$ref"] | ||
| 494 | + request_schema = schema["components"]["schemas"][request_ref.rsplit("/", 1)[1]] | ||
| 495 | + assert request_schema["required"] == ["metadata", "rawdata"] | ||
| 496 | + assert request_schema["properties"]["type"]["const"] == msg_type | ||
| 497 | + assert "application/json" in operation["post"]["responses"]["200"]["content"] | ||
| 498 | + | ||
| 499 | + invalid = demo_client.post( | ||
| 500 | + "/api/db/write", | ||
| 501 | + json=_envelope( | ||
| 502 | + "db/write", | ||
| 503 | + {"id": "record-without-value"}, | ||
| 504 | + request_id="invalid-1", | ||
| 505 | + ), | ||
| 506 | + ) | ||
| 507 | + wrong_type = demo_client.post( | ||
| 508 | + "/api/db/write", | ||
| 509 | + json=_envelope( | ||
| 510 | + "db/read", | ||
| 511 | + {"id": "record-1"}, | ||
| 512 | + request_id="invalid-2", | ||
| 513 | + ), | ||
| 514 | + ) | ||
| 515 | + | ||
| 516 | + assert invalid.status_code == 400 | ||
| 517 | + assert invalid.json()["error_code"] == "validation" | ||
| 518 | + assert wrong_type.status_code == 422 | ||
| @@ -3,12 +3,16 @@ | |||
| 3 | 3 | ||
| 4 | """SystemContext / RequestContext 单测(设计 §8):进程级/请求级、start/stop、for_request、transaction。""" | 4 | """SystemContext / RequestContext 单测(设计 §8):进程级/请求级、start/stop、for_request、transaction。""" |
| 5 | import logging | 5 | import logging |
| 6 | +from unittest.mock import AsyncMock | ||
| 6 | 7 | ||
| 7 | import fakeredis.aioredis | 8 | import fakeredis.aioredis |
| 8 | import pytest | 9 | import pytest |
| 9 | 10 | ||
| 10 | -from openjiuwen_runtime.service import RequestContext as PublicRequestContext | 11 | +from openjiuwen_runtime.service import ( |
| 11 | -from openjiuwen_runtime.service import TypedAppContext | 12 | + KubernetesUnavailable, |
| 13 | + RequestContext as PublicRequestContext, | ||
| 14 | + TypedAppContext, | ||
| 15 | +) | ||
| 12 | from openjiuwen_runtime.service.context.request_context import RequestContext | 16 | from openjiuwen_runtime.service.context.request_context import RequestContext |
| 13 | from openjiuwen_runtime.service.context.system_context import ( | 17 | from openjiuwen_runtime.service.context.system_context import ( |
| 14 | RequestContext as LegacyRequestContext, | 18 | RequestContext as LegacyRequestContext, |
| @@ -175,3 +179,50 @@ def test_from_settings_builds_redis_from_env(monkeypatch): | |||
| 175 | ctx = SystemContext.from_settings() | 179 | ctx = SystemContext.from_settings() |
| 176 | assert ctx.redis is not None | 180 | assert ctx.redis is not None |
| 177 | assert ctx._owns_redis is True # 自建 → stop 时负责关闭 | 181 | assert ctx._owns_redis is True # 自建 → stop 时负责关闭 |
| 182 | + | ||
| 183 | + | ||
| 184 | + | ||
| 185 | +async def test_kubernetes_capability_lifecycle_and_request_entry(): | ||
| 186 | + kubernetes = AsyncMock() | ||
| 187 | + kubernetes.ping.return_value = True | ||
| 188 | + ctx = SystemContext(kubernetes=kubernetes, _owns_kubernetes=True) | ||
| 189 | + | ||
| 190 | + assert ctx.require_kubernetes() is kubernetes | ||
| 191 | + request = ctx.for_request(Metadata(request_id="kubernetes-1")) | ||
| 192 | + assert request.kubernetes is kubernetes | ||
| 193 | + | ||
| 194 | + await ctx.start() | ||
| 195 | + kubernetes.start.assert_awaited_once() | ||
| 196 | + kubernetes.ping.assert_awaited_once() | ||
| 197 | + assert await ctx.readiness() == { | ||
| 198 | + "db": None, | ||
| 199 | + "redis": None, | ||
| 200 | + "kubernetes": True, | ||
| 201 | + "lock": None, | ||
| 202 | + "cache": None, | ||
| 203 | + "ready": True, | ||
| 204 | + } | ||
| 205 | + await ctx.stop() | ||
| 206 | + kubernetes.close.assert_awaited_once() | ||
| 207 | + | ||
| 208 | + | ||
| 209 | + | ||
| 210 | +async def test_external_kubernetes_capability_keeps_caller_ownership(): | ||
| 211 | + kubernetes = AsyncMock() | ||
| 212 | + kubernetes.ping.return_value = True | ||
| 213 | + ctx = SystemContext(kubernetes=kubernetes) | ||
| 214 | + | ||
| 215 | + await ctx.start() | ||
| 216 | + await ctx.stop() | ||
| 217 | + | ||
| 218 | + kubernetes.start.assert_not_awaited() | ||
| 219 | + kubernetes.close.assert_not_awaited() | ||
| 220 | + kubernetes.ping.assert_awaited_once() | ||
| 221 | + | ||
| 222 | + | ||
| 223 | + | ||
| 224 | +def test_missing_kubernetes_capability_raises_unavailable(): | ||
| 225 | + ctx = SystemContext() | ||
| 226 | + | ||
| 227 | + with pytest.raises(KubernetesUnavailable, match="not configured"): | ||
| 228 | + ctx.require_kubernetes() | ||
🟠 High Priority
变更行
system_context.py:70新增_owns_kubernetes: bool = False(默认不持有),而system_context.py:491新增的from_settings(kubernetes=...)只是把 kubernetes 透传给build_system_context;但build_system_context(bootstrap.py:195)构造SystemContext(kubernetes=kubernetes, ...)时并未传_owns_kubernetes=True,于是_owns_kubernetes恒为 False。→ 受影响行为:
SystemContext.start()中 kubernetes 分支(system_context.py:221-228)在_owns_kubernetes为 False 时跳过await self.kubernetes.start,直接执行_ping(self.kubernetes)。→ 失败模式:
FakeKubernetesOperations.ping()在start()被调用前恒返回False(self._started),KubernetesAsyncioOperations.ping()在_core_api未初始化时直接raise KubernetesUnavailable(_require_started)。因此 demo(simple_capabilities_app.py:307-318 新建FakeKubernetesOperations(...)后经build_system_context(kubernetes=...)装配、且从不自行 start)在框架启动探活时必抛KubernetesUnavailable,即使跳过启动,get_pod/create_pod/delete_pod也会因未 start 而抛KubernetesUnavailable,K8s 能力整体不可用。这是本次 diff 引入的_owns_kubernetes默认值与既有build_system_context未同步的契约断裂。建议:将
_owns_kubernetes默认值改为None(或True),并在__init__中按kubernetes is not None自动探测;同时在build_system_context/from_settings中显式传递_owns_kubernetes,保证经工厂装配的 kubernetes 会被SystemContext.start()真正start()。