已合并
feat(service): request context中实现async Redis client调用 #415
Wal1et创建于 8月11日
feat(service): request context中实现async Redis client调用 #415
已合并
共 4 个文件变更+105-1
| @@ -509,6 +509,25 @@ sequenceDiagram | |||
| 509 | - `ctx.settings` | 509 | - `ctx.settings` |
| 510 | - `ctx.logger` | 510 | - `ctx.logger` |
| 511 | 511 | ||
| 512 | +`ctx.redis` 返回 `SystemContext` 持有的共享异步 Redis 客户端,适合在框架现有 | ||
| 513 | +`kv`、`idempotency`、`queue`、`pubsub` 和 `lock` 原语无法覆盖业务需求时,直接使用 | ||
| 514 | +`redis.asyncio` 的完整命令集: | ||
| 515 | + | ||
| 516 | +```python | ||
| 517 | +async def update_score(ctx, env): | ||
| 518 | + await ctx.redis.hset( | ||
| 519 | + f"user:{ctx.user_id}", | ||
| 520 | + mapping={"score": ctx.request.score}, | ||
| 521 | + ) | ||
| 522 | + await ctx.redis.zadd("user:scores", {str(ctx.user_id): ctx.request.score}) | ||
| 523 | + return {"updated": True} | ||
| 524 | +``` | ||
| 525 | + | ||
| 526 | +该客户端由 `SystemContext` 创建并在应用停止时统一关闭;Handler 和请求清理逻辑只借用 | ||
| 527 | +它,不能调用 `close()` 或 `aclose()`。裸客户端不会自动添加服务、用户或 Bot | ||
| 528 | +命名空间,业务代码必须自行设计不会冲突的 Key。Redis 命令是异步 I/O,调用时必须 | ||
| 529 | +使用 `await`;Pipeline 中的命令先同步入队,最终的 `execute()` 必须使用 `await`。 | ||
| 530 | + | ||
| 512 | 如果需要事务,应使用框架提供的请求上下文事务入口,而不是在 Handler 内创建无法统一 | 531 | 如果需要事务,应使用框架提供的请求上下文事务入口,而不是在 Handler 内创建无法统一 |
| 513 | 管理的连接。示例中的 `DemoUserStore` 为了便于运行,使用进程内字典和 | 532 | 管理的连接。示例中的 `DemoUserStore` 为了便于运行,使用进程内字典和 |
| 514 | `asyncio.Lock`;它的数据会在进程退出后丢失,也不会在多个服务副本之间共享。 | 533 | `asyncio.Lock`;它的数据会在进程退出后丢失,也不会在多个服务副本之间共享。 |
| @@ -236,6 +236,21 @@ class RequestContext(Generic[TRequest]): | |||
| 236 | """Return the configured DB handler after checking request state.""" | 236 | """Return the configured DB handler after checking request state.""" |
| 237 | return self.require_db() | 237 | return self.require_db() |
| 238 | 238 | ||
| 239 | + | ||
| 240 | + def redis(self) -> Any: | ||
| 241 | + """Return the shared asynchronous Redis client for this request. | ||
| 242 | + | ||
| 243 | + The client is owned by :class:`SystemContext` and is only borrowed by | ||
| 244 | + the request. Callers may use the complete ``redis.asyncio`` API, but | ||
| 245 | + must not close the client from request or handler code. | ||
| 246 | + """ | ||
| 247 | + return self.require_redis() | ||
| 248 | + | ||
| 249 | + def require_redis(self) -> Any: | ||
| 250 | + """Require the shared asynchronous Redis client for an active request.""" | ||
| 251 | + self.check_interrupted() | ||
| 252 | + return self.sysctx.require_redis() | ||
| 253 | + | ||
| 239 | def require_db(self) -> Any: | 254 | def require_db(self) -> Any: |
| 240 | """Require a DB handler for this active request.""" | 255 | """Require a DB handler for this active request.""" |
| 241 | self.check_interrupted() | 256 | self.check_interrupted() |
| @@ -85,3 +85,32 @@ def test_ws_echo_returns_idx(): | |||
| 85 | data = json.loads(ws.receive_text()) | 85 | data = json.loads(ws.receive_text()) |
| 86 | assert data["ok"] is True | 86 | assert data["ok"] is True |
| 87 | assert data["rawdata"] == {"echo": "hi", "idx": 1} | 87 | assert data["rawdata"] == {"echo": "hi", "idx": 1} |
| 88 | + | ||
| 89 | + | ||
| 90 | + | ||
| 91 | +async def test_rest_handler_can_use_raw_async_redis_commands(): | ||
| 92 | + redis = fakeredis.aioredis.FakeRedis() | ||
| 93 | + app = App(lambda: SystemContext(redis=redis)) | ||
| 94 | + | ||
| 95 | + | ||
| 96 | + async def save_profile(ctx, env: Envelope): | ||
| 97 | + key = f"user:{ctx.user_id}" | ||
| 98 | + await ctx.redis.hset(key, mapping=env.rawdata) | ||
| 99 | + name = await ctx.redis.hget(key, "name") | ||
| 100 | + return {"name": name.decode("utf-8")} | ||
| 101 | + | ||
| 102 | + async with httpx.AsyncClient( | ||
| 103 | + transport=httpx.ASGITransport(app=app.asgi), | ||
| 104 | + base_url="http://test", | ||
| 105 | + ) as client: | ||
| 106 | + response = await client.post( | ||
| 107 | + "/api/redis.profile", | ||
| 108 | + json={ | ||
| 109 | + "type": "redis.profile", | ||
| 110 | + "metadata": {"request_id": "redis-1", "user_id": "user-1"}, | ||
| 111 | + "rawdata": {"name": "Alice"}, | ||
| 112 | + }, | ||
| 113 | + ) | ||
| 114 | + | ||
| 115 | + assert response.status_code == 200 | ||
| 116 | + assert response.json()["rawdata"] == {"name": "Alice"} | ||
| @@ -1,7 +1,7 @@ | |||
| 1 | # coding: utf-8 | 1 | # coding: utf-8 |
| 2 | # Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved | 2 | # Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved |
| 3 | 3 | ||
| 4 | -"""DB, transaction, audit, and lock capabilities on RequestContext.""" | 4 | +"""DB, Redis, transaction, audit, and lock capabilities on RequestContext.""" |
| 5 | from __future__ import annotations | 5 | from __future__ import annotations |
| 6 | 6 | ||
| 7 | import logging | 7 | import logging |
| @@ -166,6 +166,47 @@ async def test_db_helpers_check_request_state_and_missing_database(): | |||
| 166 | await interrupted.db_count("users") | 166 | await interrupted.db_count("users") |
| 167 | 167 | ||
| 168 | 168 | ||
| 169 | + | ||
| 170 | +async def test_request_context_exposes_shared_async_redis_client(mocker): | ||
| 171 | + redis = fakeredis.aioredis.FakeRedis() | ||
| 172 | + close_spy = mocker.spy(redis, "aclose") | ||
| 173 | + sysctx = SystemContext(redis=redis) | ||
| 174 | + first = sysctx.for_request(_env()) | ||
| 175 | + second = sysctx.for_request(_env()) | ||
| 176 | + | ||
| 177 | + assert first.redis is redis | ||
| 178 | + assert first.require_redis() is redis | ||
| 179 | + assert second.redis is redis | ||
| 180 | + | ||
| 181 | + await first.redis.hset("user:1", mapping={"name": "Alice"}) | ||
| 182 | + await first.redis.zadd("user:scores", {"user:1": 10}) | ||
| 183 | + | ||
| 184 | + assert await second.redis.hget("user:1", "name") == b"Alice" | ||
| 185 | + assert await second.redis.zscore("user:scores", "user:1") == 10 | ||
| 186 | + | ||
| 187 | + await first.close() | ||
| 188 | + close_spy.assert_not_awaited() | ||
| 189 | + assert await redis.ping() is True | ||
| 190 | + | ||
| 191 | + | ||
| 192 | + | ||
| 193 | +def test_request_context_redis_requires_configuration_and_active_request(): | ||
| 194 | + missing = SystemContext().for_request(_env()) | ||
| 195 | + with pytest.raises(RedisUnavailable): | ||
| 196 | + _ = missing.redis | ||
| 197 | + with pytest.raises(RedisUnavailable): | ||
| 198 | + missing.require_redis() | ||
| 199 | + | ||
| 200 | + redis = fakeredis.aioredis.FakeRedis() | ||
| 201 | + interrupted = SystemContext(redis=redis).for_request(_env()) | ||
| 202 | + interrupted.interrupt("stop redis work") | ||
| 203 | + | ||
| 204 | + with pytest.raises(Interrupted, match="stop redis work"): | ||
| 205 | + _ = interrupted.redis | ||
| 206 | + with pytest.raises(Interrupted, match="stop redis work"): | ||
| 207 | + interrupted.require_redis() | ||
| 208 | + | ||
| 209 | + | ||
| 169 | 210 | ||
| 170 | async def test_transaction_yields_one_session_and_db_helpers_stay_independent(): | 211 | async def test_transaction_yields_one_session_and_db_helpers_stay_independent(): |
| 171 | db = _FakeDb() | 212 | db = _FakeDb() |