已合并
feat(service):上下文能力demo、集成Kubernetes至上下文 #418
feat(service):上下文能力demo、集成Kubernetes至上下文 #418
已合并
m0u55e创建于 8月14日
共 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+@dataclass(frozen=True)
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+ @classmethod
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+ @field_serializer("updated_at")
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+ @application.handle(
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+ @application.handle(
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+ @application.handle(
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+ @application.handle(
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+ @application.handle(
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+ @application.handle(
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+ @application.handle(
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+ @application.handle(
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 # locks139 # 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+)
32from .request_context import RequestContext, TypedAppContext40from .request_context import RequestContext, TypedAppContext
33from .system_context import SystemContext41from .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+ @classmethod
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+ @staticmethod
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+ @staticmethod
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+@dataclass(frozen=True)
13+class PodCreateSpec:
14+ """Restricted input used to create one managed Pod."""
15+ 
16+ name: str
17+ image: str
18+ 
19+ 
20+@dataclass(frozen=True)
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+@dataclass(frozen=True)
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+@runtime_checkable
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 
25if TYPE_CHECKING:25if TYPE_CHECKING:
26 from .cache import Cache26 from .cache import Cache
27+ from .kubernetes import KubernetesOperations
27 from .locks import LockManager28 from .locks import LockManager
28 from .system_context import SystemContext29 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+ @property
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
36from .audit import AuditEvent, AuditLogger, LoggingAuditLogger, NoopAuditLogger38from .audit import AuditEvent, AuditLogger, LoggingAuditLogger, NoopAuditLogger
37from .request_context import RequestContext39from .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,
atomgit-bot
atomgit-botatomgit-bot8月14日

🟠 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()。

likedislike
不准确?
67 ) -> None:71 ) -> None:
68 self.redis = redis72 self.redis = redis
69 self.db = db73 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 cache94 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_redis98 self._owns_redis = _owns_redis
@@ -102,6 +107,7 @@ class SystemContext:
102 if _owns_cache_backend is None107 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 = False111 self._started = False
106 self._stopped = False112 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_backend148 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 = db157 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 = backend171 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 continue302 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 return325 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", None335 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 
33class FrameworkError(Exception):35class FrameworkError(Exception):
@@ -125,6 +127,18 @@ class CacheUnavailable(FrameworkError):
125 code = ErrorCode.CACHE_UNAVAILABLE127 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@runtime_checkable142@runtime_checkable
129class _HasCode(Protocol):143class _HasCode(Protocol):
130 code: str144 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 = [
25etcd = [25etcd = [
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]
30default-groups = ['dev']33default-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_UNAVAILABLE45 assert DatabaseUnavailable("x").code == ErrorCode.DATABASE_UNAVAILABLE
44 assert RedisUnavailable("x").code == ErrorCode.REDIS_UNAVAILABLE46 assert RedisUnavailable("x").code == ErrorCode.REDIS_UNAVAILABLE
45 assert CacheUnavailable("x").code == ErrorCode.CACHE_UNAVAILABLE47 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) == 50386 assert http_status_for(ErrorCode.DATABASE_UNAVAILABLE) == 503
83 assert http_status_for(ErrorCode.REDIS_UNAVAILABLE) == 50387 assert http_status_for(ErrorCode.REDIS_UNAVAILABLE) == 503
84 assert http_status_for(ErrorCode.CACHE_UNAVAILABLE) == 50388 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) == 50091 assert http_status_for(ErrorCode.INTERNAL) == 500
86 # 未知 code → 500(fail-safe)92 # 未知 code → 500(fail-safe)
87 assert http_status_for("totally-unknown") == 50093 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+@pytest.mark.unit
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+@pytest.mark.unit
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+@pytest.mark.unit
108+@pytest.mark.parametrize(
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+@pytest.mark.unit
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+@pytest.mark.unit
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+@pytest.mark.unit
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+@pytest.mark.unit
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+@pytest.mark.unit
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+@pytest.mark.unit
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+@dataclass(frozen=True)
31+class _InvalidRedisConfigCase:
32+ mode: str
33+ ttl: str
34+ redis_url: str
35+ message: str
36+ 
37+ 
38+@pytest.fixture(scope="module")
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+@pytest.fixture
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+@pytest.mark.unit
105+@pytest.mark.parametrize(
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+@pytest.mark.unit
132+@pytest.mark.parametrize(
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+@pytest.mark.unit
169+@pytest.mark.parametrize(
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+@pytest.mark.unit
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+@pytest.mark.unit
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+@pytest.mark.unit
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+@pytest.mark.unit
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+@pytest.mark.unit
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+@pytest.mark.unit
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。"""
5import logging5import logging
6+from unittest.mock import AsyncMock
6 7 
7import fakeredis.aioredis8import fakeredis.aioredis
8import pytest9import pytest
9 10 
10-from openjiuwen_runtime.service import RequestContext as PublicRequestContext11+from openjiuwen_runtime.service import (
11-from openjiuwen_runtime.service import TypedAppContext12+ KubernetesUnavailable,
13+ RequestContext as PublicRequestContext,
14+ TypedAppContext,
15+)
12from openjiuwen_runtime.service.context.request_context import RequestContext16from openjiuwen_runtime.service.context.request_context import RequestContext
13from openjiuwen_runtime.service.context.system_context import (17from 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 None180 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+@pytest.mark.unit
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+@pytest.mark.unit
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+@pytest.mark.unit
224+def test_missing_kubernetes_capability_raises_unavailable():
225+ ctx = SystemContext()
226+ 
227+ with pytest.raises(KubernetesUnavailable, match="not configured"):
228+ ctx.require_kubernetes()