已合并
feat(service): request context中实现async Redis client调用 #415
feat(service): request context中实现async Redis client调用 #415
已合并
Wal1et创建于 8月11日
共 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+ @property
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 True86 assert data["ok"] is True
87 assert data["rawdata"] == {"echo": "hi", "idx": 1}87 assert data["rawdata"] == {"echo": "hi", "idx": 1}
88+ 
89+ 
90+@pytest.mark.system
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+ @app.handle("redis.profile")
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-81# coding: utf-8
2# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved2# 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."""
5from __future__ import annotations5from __future__ import annotations
6 6 
7import logging7import 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+@pytest.mark.unit
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+@pytest.mark.unit
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@pytest.mark.unit210@pytest.mark.unit
170async def test_transaction_yields_one_session_and_db_helpers_stay_independent():211async def test_transaction_yields_one_session_and_db_helpers_stay_independent():
171 db = _FakeDb()212 db = _FakeDb()