Files
insistence 5cd6d4788a feat: 新增多worker运行支持 (#76)
* feat: 定时任务新增redis分布式锁以在多worker下正常运行

* refactor: 重构日志系统以支持多worker

* perf: 增强scheduler在多worker下的调度机制

* perf: 增强锁丢失调度逻辑与异常容错

* perf: 优化trace中间件和日志打印显示

* chore: 更新.gitignore规则

* perf: 优化database的engine和session创建方法

* perf: 优化scheduler内部的sqlalchemy日志显示

* perf: 优化scheduler和应用生命周期

* fix: 修复重新获取锁时重复调用的问题

* perf: 优化函数无用参数

* perf: 优化env配置

* fix: 修复lint错误
2026-02-06 09:07:46 +08:00

55 lines
1.8 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
from collections.abc import AsyncGenerator
from contextlib import asynccontextmanager
from starlette.types import Message, Scope
from .ctx import TraceCtx
class Span:
"""
整个http生命周期:
request(before) --> request(after) --> response(before) --> response(after)
"""
def __init__(self, scope: Scope) -> None:
self.scope = scope
async def request_before(self) -> None:
"""
request_before: 处理header信息等, 如记录请求体信息
"""
TraceCtx.set_trace_id()
TraceCtx.set_request_id()
TraceCtx.set_span_id()
TraceCtx.set_request_path(self.scope.get('path', ''))
TraceCtx.set_request_method(self.scope.get('method', ''))
async def request_after(self, message: Message) -> Message:
"""
request_after: 处理请求bytes, 如记录请求参数
example:
message: {'type': 'http.request', 'body': b'{\r\n "name": "\xe8\x8b\x8f\xe8\x8b\x8f\xe8\x8b\x8f"\r\n}', 'more_body': False}
"""
return message
async def response(self, message: Message) -> Message:
"""
if message['type'] == "http.response.start": -----> request-before
pass
if message['type'] == "http.response.body": -----> request-after
message.get('body', b'')
pass
"""
if message['type'] == 'http.response.start':
message['headers'].append((b'request-id', TraceCtx.get_request_id().encode()))
message['headers'].append((b'trace-id', TraceCtx.get_trace_id().encode()))
message['headers'].append((b'span-id', TraceCtx.get_span_id().encode()))
return message
@asynccontextmanager
async def get_current_span(scope: Scope) -> AsyncGenerator[Span, None]:
yield Span(scope)