feat: 请求日志和 usage 记录中追踪代理节点信息

- 新增 resolve_proxy_info/get_proxy_label 提取脱敏的代理摘要
- StreamContext 增加 proxy_info 字段,流式/非流式全路径写入 metadata
- ProxyNode 缓存补充 name 字段,日志输出代理节点标识
- aether-proxy 连接和转发日志从 debug 提升为 info 级别
- TaskService 代理不可用日志补充 node_id 细节
This commit is contained in:
fawney19
2026-02-08 00:49:43 +08:00
parent 5384ffd403
commit 4a1029961a
9 changed files with 153 additions and 27 deletions

View File

@@ -4,7 +4,7 @@ use std::sync::Arc;
use hyper::body::Incoming; use hyper::body::Incoming;
use hyper::{Request, Response}; use hyper::{Request, Response};
use tokio::net::TcpStream; use tokio::net::TcpStream;
use tracing::{debug, warn}; use tracing::{debug, info, warn};
use crate::auth; use crate::auth;
use crate::config::Config; use crate::config::Config;
@@ -53,7 +53,7 @@ pub async fn handle_connect(
} }
}; };
debug!(target = %target_addr, "CONNECT tunnel establishing"); info!(target = %target_addr, "CONNECT tunnel establishing");
// Connect to target // Connect to target
let target_stream = match TcpStream::connect(target_addr).await { let target_stream = match TcpStream::connect(target_addr).await {
@@ -74,7 +74,7 @@ pub async fn handle_connect(
match tokio::io::copy_bidirectional(&mut upgraded, &mut target).await { match tokio::io::copy_bidirectional(&mut upgraded, &mut target).await {
Ok((from_client, from_target)) => { Ok((from_client, from_target)) => {
debug!( info!(
from_client, from_client,
from_target, from_target,
"CONNECT tunnel closed" "CONNECT tunnel closed"

View File

@@ -4,7 +4,7 @@ use std::sync::Arc;
use http_body_util::{BodyExt, Full}; use http_body_util::{BodyExt, Full};
use hyper::body::Incoming; use hyper::body::Incoming;
use hyper::{Request, Response}; use hyper::{Request, Response};
use tracing::{debug, warn}; use tracing::{debug, info, warn};
use crate::auth; use crate::auth;
use crate::config::Config; use crate::config::Config;
@@ -56,7 +56,7 @@ pub async fn handle_plain(
} }
}; };
debug!(target = %target_addr, method = %req.method(), "HTTP proxy forwarding"); info!(target = %target_addr, method = %req.method(), "HTTP proxy forwarding");
// Build outgoing request (strip proxy headers, use relative URI) // Build outgoing request (strip proxy headers, use relative URI)
let path_and_query = uri let path_and_query = uri
@@ -116,6 +116,7 @@ pub async fn handle_plain(
match sender.send_request(outgoing).await { match sender.send_request(outgoing).await {
Ok(resp) => { Ok(resp) => {
info!(target = %target_addr, status = resp.status().as_u16(), "HTTP proxy response");
// Stream the response body directly — no buffering // Stream the response body directly — no buffering
let (parts, body) = resp.into_parts(); let (parts, body) = resp.into_parts();
let body: BoxBody = body let body: BoxBody = body

View File

@@ -53,7 +53,7 @@ pub async fn run(
} }
}; };
debug!(peer = %peer_addr, "new connection"); info!(peer = %peer_addr, "new connection");
let config = Arc::clone(&config); let config = Arc::clone(&config);
let node_id = Arc::clone(&node_id); let node_id = Arc::clone(&node_id);

View File

@@ -913,9 +913,15 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
# Capture the selected base_url from transport (used by some envelopes for failover). # Capture the selected base_url from transport (used by some envelopes for failover).
ctx.selected_base_url = envelope.capture_selected_base_url() if envelope else None ctx.selected_base_url = envelope.capture_selected_base_url() if envelope else None
# 记录代理信息
from src.clients.http_client import get_proxy_label, resolve_proxy_info
ctx.proxy_info = resolve_proxy_info(provider.proxy)
proxy_label = get_proxy_label(ctx.proxy_info)
logger.debug( logger.debug(
f" [{self.request_id}] 发送流式请求: Provider={provider.name}, " f" [{self.request_id}] 发送流式请求: Provider={provider.name}, "
f"模型={ctx.model} -> {mapped_model or '无映射'}" f"模型={ctx.model} -> {mapped_model or '无映射'}, 代理={proxy_label}"
) )
# If upstream is forced to non-stream mode, we execute a sync request and then # If upstream is forced to non-stream mode, we execute a sync request and then
@@ -1307,6 +1313,10 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
# 失败时返回给客户端的是 JSON 错误响应 # 失败时返回给客户端的是 JSON 错误响应
client_response_headers = {"content-type": "application/json"} client_response_headers = {"content-type": "application/json"}
stream_fail_metadata: dict[str, Any] | None = None
if ctx.proxy_info:
stream_fail_metadata = {"proxy": ctx.proxy_info}
await self.telemetry.record_failure( await self.telemetry.record_failure(
provider=ctx.provider_name or "unknown", provider=ctx.provider_name or "unknown",
model=ctx.model, model=ctx.model,
@@ -1324,6 +1334,7 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
endpoint_api_format=ctx.provider_api_format or None, endpoint_api_format=ctx.provider_api_format or None,
has_format_conversion=ctx.has_format_conversion, has_format_conversion=ctx.has_format_conversion,
target_model=ctx.mapped_model, target_model=ctx.mapped_model,
request_metadata=stream_fail_metadata,
) )
# ==================== 非流式处理 ==================== # ==================== 非流式处理 ====================
@@ -1372,6 +1383,7 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
endpoint_id: str | None = None # Endpoint ID用于失败记录 endpoint_id: str | None = None # Endpoint ID用于失败记录
key_id: str | None = None # Key ID用于失败记录 key_id: str | None = None # Key ID用于失败记录
mapped_model_result: str | None = None # 映射后的目标模型名(用于 Usage 记录) mapped_model_result: str | None = None # 映射后的目标模型名(用于 Usage 记录)
sync_proxy_info: dict[str, Any] | None = None # 代理信息(用于 Usage 记录)
async def sync_request_func( async def sync_request_func(
provider: Provider, provider: Provider,
@@ -1382,6 +1394,7 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
nonlocal provider_name, response_json, status_code, response_headers nonlocal provider_name, response_json, status_code, response_headers
nonlocal provider_request_headers, provider_request_body, mapped_model_result nonlocal provider_request_headers, provider_request_body, mapped_model_result
nonlocal provider_api_format_for_error, client_api_format_for_error, needs_conversion_for_error nonlocal provider_api_format_for_error, client_api_format_for_error, needs_conversion_for_error
nonlocal sync_proxy_info
provider_name = str(provider.name) provider_name = str(provider.name)
provider_api_format = str(endpoint.api_format or api_format) provider_api_format = str(endpoint.api_format or api_format)
@@ -1533,9 +1546,16 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
# 非流式:必须在 build_provider_url 调用后立即缓存(避免 contextvar 被后续调用覆盖) # 非流式:必须在 build_provider_url 调用后立即缓存(避免 contextvar 被后续调用覆盖)
selected_base_url_cached = envelope.capture_selected_base_url() if envelope else None selected_base_url_cached = envelope.capture_selected_base_url() if envelope else None
# 记录代理信息
from src.clients.http_client import get_proxy_label, resolve_proxy_info
sync_proxy_info = resolve_proxy_info(provider.proxy)
_proxy_label = get_proxy_label(sync_proxy_info)
logger.info( logger.info(
f" [{self.request_id}] 发送{'上游流式(聚合)' if upstream_is_stream else '非流式'}请求: " f" [{self.request_id}] 发送{'上游流式(聚合)' if upstream_is_stream else '非流式'}请求: "
f"Provider={provider.name}, 模型={model} -> {mapped_model or '无映射'}" f"Provider={provider.name}, 模型={model} -> {mapped_model or '无映射'}, "
f"代理={_proxy_label}"
) )
logger.debug(f" [{self.request_id}] 请求URL: {redact_url_for_log(url)}") logger.debug(f" [{self.request_id}] 请求URL: {redact_url_for_log(url)}")
@@ -1775,7 +1795,9 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
client_response_headers = filter_proxy_response_headers(response_headers) client_response_headers = filter_proxy_response_headers(response_headers)
client_response_headers["content-type"] = "application/json" client_response_headers["content-type"] = "application/json"
request_metadata = self._build_request_metadata() request_metadata = self._build_request_metadata() or {}
if sync_proxy_info:
request_metadata["proxy"] = sync_proxy_info
total_cost = await self.telemetry.record_success( total_cost = await self.telemetry.record_success(
provider=provider_name, provider=provider_name,
model=model, model=model,
@@ -1803,7 +1825,7 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
provider_api_key_id=key_id, provider_api_key_id=key_id,
# 模型映射信息 # 模型映射信息
target_model=mapped_model_result, target_model=mapped_model_result,
request_metadata=request_metadata, request_metadata=request_metadata or None,
) )
logger.debug(f"{self.FORMAT_ID} 非流式响应完成") logger.debug(f"{self.FORMAT_ID} 非流式响应完成")
@@ -1826,7 +1848,9 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
# 记录实际发送给 Provider 的请求体,便于排查问题根因 # 记录实际发送给 Provider 的请求体,便于排查问题根因
response_time_ms = self.elapsed_ms() response_time_ms = self.elapsed_ms()
actual_request_body = provider_request_body or original_request_body actual_request_body = provider_request_body or original_request_body
request_metadata = self._build_request_metadata() request_metadata = self._build_request_metadata() or {}
if sync_proxy_info:
request_metadata["proxy"] = sync_proxy_info
await self.telemetry.record_failure( await self.telemetry.record_failure(
provider=provider_name or "unknown", provider=provider_name or "unknown",
model=model, model=model,
@@ -1836,7 +1860,7 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
request_body=actual_request_body, request_body=actual_request_body,
error_message=str(e), error_message=str(e),
is_stream=False, is_stream=False,
request_metadata=request_metadata, request_metadata=request_metadata or None,
) )
client_format = (client_api_format_for_error or "").upper() client_format = (client_api_format_for_error or "").upper()
provider_format = (provider_api_format_for_error or client_format).upper() provider_format = (provider_api_format_for_error or client_format).upper()
@@ -1851,7 +1875,9 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
except UpstreamClientException as e: except UpstreamClientException as e:
response_time_ms = self.elapsed_ms() response_time_ms = self.elapsed_ms()
actual_request_body = provider_request_body or original_request_body actual_request_body = provider_request_body or original_request_body
request_metadata = self._build_request_metadata() request_metadata = self._build_request_metadata() or {}
if sync_proxy_info:
request_metadata["proxy"] = sync_proxy_info
await self.telemetry.record_failure( await self.telemetry.record_failure(
provider=provider_name or "unknown", provider=provider_name or "unknown",
model=model, model=model,
@@ -1903,7 +1929,9 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
elif isinstance(e, httpx.HTTPStatusError) and hasattr(e, "response"): elif isinstance(e, httpx.HTTPStatusError) and hasattr(e, "response"):
error_response_headers = dict(e.response.headers) error_response_headers = dict(e.response.headers)
request_metadata = self._build_request_metadata() request_metadata = self._build_request_metadata() or {}
if sync_proxy_info:
request_metadata["proxy"] = sync_proxy_info
await self.telemetry.record_failure( await self.telemetry.record_failure(
provider=provider_name or "unknown", provider=provider_name or "unknown",
model=model, model=model,

View File

@@ -890,6 +890,12 @@ class CliMessageHandlerBase(BaseMessageHandler):
# Capture the selected base_url from transport (used by some envelopes for failover). # Capture the selected base_url from transport (used by some envelopes for failover).
ctx.selected_base_url = envelope.capture_selected_base_url() if envelope else None ctx.selected_base_url = envelope.capture_selected_base_url() if envelope else None
# 记录代理信息sync-bridge 路径,早于流式路径执行)
from src.clients.http_client import get_proxy_label as _gpl
from src.clients.http_client import resolve_proxy_info as _rpi
ctx.proxy_info = _rpi(provider.proxy)
# If upstream is forced to non-stream mode, we execute a sync request and then # If upstream is forced to non-stream mode, we execute a sync request and then
# simulate streaming to the client (sync -> stream bridge). # simulate streaming to the client (sync -> stream bridge).
if not upstream_is_stream: if not upstream_is_stream:
@@ -1073,12 +1079,14 @@ class CliMessageHandlerBase(BaseMessageHandler):
# 优先使用 Provider 配置,否则使用全局配置 # 优先使用 Provider 配置,否则使用全局配置
request_timeout = provider.stream_first_byte_timeout or config.stream_first_byte_timeout request_timeout = provider.stream_first_byte_timeout or config.stream_first_byte_timeout
_proxy_label = _gpl(ctx.proxy_info)
logger.debug( logger.debug(
f" └─ [{self.request_id}] 发送流式请求: " f" └─ [{self.request_id}] 发送流式请求: "
f"Provider={provider.name}, Endpoint={endpoint.id[:8] if endpoint.id else 'N/A'}..., " f"Provider={provider.name}, Endpoint={endpoint.id[:8] if endpoint.id else 'N/A'}..., "
f"Key=***{key.api_key[-4:] if key.api_key else 'N/A'}, " f"Key=***{key.api_key[-4:] if key.api_key else 'N/A'}, "
f"原始模型={ctx.model}, 映射后={mapped_model or '无映射'}, URL模型={url_model}, " f"原始模型={ctx.model}, 映射后={mapped_model or '无映射'}, URL模型={url_model}, "
f"timeout={request_timeout}s" f"timeout={request_timeout}s, 代理={_proxy_label}"
) )
# 创建 HTTP 客户端(支持代理配置,从 Provider 读取) # 创建 HTTP 客户端(支持代理配置,从 Provider 读取)
@@ -2802,6 +2810,7 @@ class CliMessageHandlerBase(BaseMessageHandler):
mapped_model_result = None # 映射后的目标模型名(用于 Usage 记录) mapped_model_result = None # 映射后的目标模型名(用于 Usage 记录)
response_metadata_result: dict[str, Any] = {} # Provider 响应元数据 response_metadata_result: dict[str, Any] = {} # Provider 响应元数据
needs_conversion = False # 是否需要格式转换(由 candidate 决定) needs_conversion = False # 是否需要格式转换(由 candidate 决定)
sync_proxy_info: dict[str, Any] | None = None # 代理信息
# 可变请求体容器:允许 TaskService 在遇到 Thinking 签名错误时整流请求体后重试 # 可变请求体容器:允许 TaskService 在遇到 Thinking 签名错误时整流请求体后重试
# 结构: {"body": 实际请求体, "_rectified": 是否已整流, "_rectified_this_turn": 本轮是否整流} # 结构: {"body": 实际请求体, "_rectified": 是否已整流, "_rectified_this_turn": 本轮是否整流}
@@ -2813,7 +2822,7 @@ class CliMessageHandlerBase(BaseMessageHandler):
key: ProviderAPIKey, key: ProviderAPIKey,
candidate: ProviderCandidate, candidate: ProviderCandidate,
) -> dict[str, Any]: ) -> dict[str, Any]:
nonlocal provider_name, response_json, status_code, response_headers, provider_api_format, provider_request_headers, provider_request_body, mapped_model_result, response_metadata_result, needs_conversion nonlocal provider_name, response_json, status_code, response_headers, provider_api_format, provider_request_headers, provider_request_body, mapped_model_result, response_metadata_result, needs_conversion, sync_proxy_info
provider_name = str(provider.name) provider_name = str(provider.name)
provider_api_format = str(endpoint.api_format) if endpoint.api_format else "" provider_api_format = str(endpoint.api_format) if endpoint.api_format else ""
@@ -2948,11 +2957,18 @@ class CliMessageHandlerBase(BaseMessageHandler):
# 非流式:必须在 build_provider_url 调用后立即缓存(避免 contextvar 被后续调用覆盖) # 非流式:必须在 build_provider_url 调用后立即缓存(避免 contextvar 被后续调用覆盖)
selected_base_url_cached = envelope.capture_selected_base_url() if envelope else None selected_base_url_cached = envelope.capture_selected_base_url() if envelope else None
# 记录代理信息
from src.clients.http_client import get_proxy_label, resolve_proxy_info
sync_proxy_info = resolve_proxy_info(provider.proxy)
_proxy_label = get_proxy_label(sync_proxy_info)
logger.info( logger.info(
f" └─ [{self.request_id}] 发送{'上游流式(聚合)' if upstream_is_stream else '非流式'}请求: " f" └─ [{self.request_id}] 发送{'上游流式(聚合)' if upstream_is_stream else '非流式'}请求: "
f"Provider={provider.name}, Endpoint={endpoint.id[:8] if endpoint.id else 'N/A'}..., " f"Provider={provider.name}, Endpoint={endpoint.id[:8] if endpoint.id else 'N/A'}..., "
f"Key=***{key.api_key[-4:] if key.api_key else 'N/A'}, " f"Key=***{key.api_key[-4:] if key.api_key else 'N/A'}, "
f"原始模型={model}, 映射后={mapped_model or '无映射'}, URL模型={url_model}" f"原始模型={model}, 映射后={mapped_model or '无映射'}, URL模型={url_model}, "
f"代理={_proxy_label}"
) )
# 获取复用的 HTTP 客户端(支持代理配置,从 Provider 读取) # 获取复用的 HTTP 客户端(支持代理配置,从 Provider 读取)
@@ -3190,7 +3206,9 @@ class CliMessageHandlerBase(BaseMessageHandler):
client_response_headers = filter_proxy_response_headers(response_headers) client_response_headers = filter_proxy_response_headers(response_headers)
client_response_headers["content-type"] = "application/json" client_response_headers["content-type"] = "application/json"
request_metadata = self._build_request_metadata() request_metadata = self._build_request_metadata() or {}
if sync_proxy_info:
request_metadata["proxy"] = sync_proxy_info
total_cost = await self.telemetry.record_success( total_cost = await self.telemetry.record_success(
provider=provider_name, provider=provider_name,
model=model, model=model,
@@ -3219,7 +3237,7 @@ class CliMessageHandlerBase(BaseMessageHandler):
target_model=mapped_model_result, target_model=mapped_model_result,
# Provider 响应元数据(如 Gemini 的 modelVersion # Provider 响应元数据(如 Gemini 的 modelVersion
response_metadata=response_metadata_result if response_metadata_result else None, response_metadata=response_metadata_result if response_metadata_result else None,
request_metadata=request_metadata, request_metadata=request_metadata or None,
) )
logger.info(f"{self.FORMAT_ID} 非流式响应处理完成") logger.info(f"{self.FORMAT_ID} 非流式响应处理完成")
@@ -3236,7 +3254,9 @@ class CliMessageHandlerBase(BaseMessageHandler):
# 记录实际发送给 Provider 的请求体,便于排查问题根因 # 记录实际发送给 Provider 的请求体,便于排查问题根因
response_time_ms = int((time.time() - sync_start_time) * 1000) response_time_ms = int((time.time() - sync_start_time) * 1000)
actual_request_body = provider_request_body or original_request_body actual_request_body = provider_request_body or original_request_body
request_metadata = self._build_request_metadata() request_metadata = self._build_request_metadata() or {}
if sync_proxy_info:
request_metadata["proxy"] = sync_proxy_info
await self.telemetry.record_failure( await self.telemetry.record_failure(
provider=provider_name or "unknown", provider=provider_name or "unknown",
model=model, model=model,
@@ -3247,7 +3267,7 @@ class CliMessageHandlerBase(BaseMessageHandler):
error_message=str(e), error_message=str(e),
is_stream=False, is_stream=False,
api_format=api_format, api_format=api_format,
request_metadata=request_metadata, request_metadata=request_metadata or None,
) )
raise raise
@@ -3272,7 +3292,9 @@ class CliMessageHandlerBase(BaseMessageHandler):
elif isinstance(e, httpx.HTTPStatusError) and hasattr(e, "response"): elif isinstance(e, httpx.HTTPStatusError) and hasattr(e, "response"):
error_response_headers = dict(e.response.headers) error_response_headers = dict(e.response.headers)
request_metadata = self._build_request_metadata() request_metadata = self._build_request_metadata() or {}
if sync_proxy_info:
request_metadata["proxy"] = sync_proxy_info
await self.telemetry.record_failure( await self.telemetry.record_failure(
provider=provider_name or "unknown", provider=provider_name or "unknown",
model=model, model=model,
@@ -3292,7 +3314,7 @@ class CliMessageHandlerBase(BaseMessageHandler):
has_format_conversion=is_format_converted(provider_api_format, str(api_format)), has_format_conversion=is_format_converted(provider_api_format, str(api_format)),
# 模型映射信息 # 模型映射信息
target_model=mapped_model_result, target_model=mapped_model_result,
request_metadata=request_metadata, request_metadata=request_metadata or None,
) )
raise raise

View File

@@ -112,6 +112,9 @@ class StreamContext:
perf_sampled: bool = False perf_sampled: bool = False
perf_metrics: dict[str, Any] = field(default_factory=dict) perf_metrics: dict[str, Any] = field(default_factory=dict)
# 代理信息(用于 usage 记录和日志)
proxy_info: dict[str, Any] | None = None
# 流式格式转换状态(跨 chunk 追踪) # 流式格式转换状态(跨 chunk 追踪)
stream_conversion_state: StreamState | None = None stream_conversion_state: StreamState | None = None
@@ -142,6 +145,7 @@ class StreamContext:
self.response_id = None self.response_id = None
self.final_usage = None self.final_usage = None
self.final_response = None self.final_response = None
self.proxy_info = None
self.stream_conversion_state = None self.stream_conversion_state = None
self.needs_conversion = False self.needs_conversion = False
self.selected_base_url = None self.selected_base_url = None

View File

@@ -201,6 +201,8 @@ class StreamTelemetryRecorder:
metadata: dict[str, Any] = {"stream": True, "content_length": ctx.data_count} metadata: dict[str, Any] = {"stream": True, "content_length": ctx.data_count}
if ctx.perf_metrics: if ctx.perf_metrics:
metadata["perf"] = ctx.perf_metrics metadata["perf"] = ctx.perf_metrics
if ctx.proxy_info:
metadata["proxy"] = ctx.proxy_info
await writer.record_success( await writer.record_success(
provider=ctx.provider_name or "unknown", provider=ctx.provider_name or "unknown",
@@ -250,6 +252,8 @@ class StreamTelemetryRecorder:
metadata: dict[str, Any] = {"stream": True, "content_length": ctx.data_count} metadata: dict[str, Any] = {"stream": True, "content_length": ctx.data_count}
if ctx.perf_metrics: if ctx.perf_metrics:
metadata["perf"] = ctx.perf_metrics metadata["perf"] = ctx.perf_metrics
if ctx.proxy_info:
metadata["proxy"] = ctx.proxy_info
await writer.record_failure( await writer.record_failure(
provider=ctx.provider_name or "unknown", provider=ctx.provider_name or "unknown",
@@ -297,6 +301,8 @@ class StreamTelemetryRecorder:
metadata: dict[str, Any] = {"stream": True, "content_length": ctx.data_count} metadata: dict[str, Any] = {"stream": True, "content_length": ctx.data_count}
if ctx.perf_metrics: if ctx.perf_metrics:
metadata["perf"] = ctx.perf_metrics metadata["perf"] = ctx.perf_metrics
if ctx.proxy_info:
metadata["proxy"] = ctx.proxy_info
await writer.record_cancelled( await writer.record_cancelled(
provider=ctx.provider_name or "unknown", provider=ctx.provider_name or "unknown",

View File

@@ -40,8 +40,8 @@ def _get_proxy_node_info(node_id: str) -> dict[str, Any] | None:
读取 ProxyNode 信息(带内存 TTL 缓存) 读取 ProxyNode 信息(带内存 TTL 缓存)
Returns: Returns:
aether-proxy 节点: {"ip": str, "port": int} aether-proxy 节点: {"ip": str, "port": int, "name": str, ...}
手动节点: {"is_manual": True, "proxy_url": str, "username": str|None, "password": str|None} 手动节点: {"is_manual": True, "name": str, "proxy_url": str, ...}
不存在/非在线: None 不存在/非在线: None
""" """
now = time.time() now = time.time()
@@ -68,12 +68,14 @@ def _get_proxy_node_info(node_id: str) -> dict[str, Any] | None:
if node.is_manual: if node.is_manual:
value: dict[str, Any] = { value: dict[str, Any] = {
"is_manual": True, "is_manual": True,
"name": node.name,
"proxy_url": node.proxy_url, "proxy_url": node.proxy_url,
"username": node.proxy_username, "username": node.proxy_username,
"password": node.proxy_password, "password": node.proxy_password,
} }
else: else:
value = { value = {
"name": node.name,
"ip": node.ip, "ip": node.ip,
"port": node.port, "port": node.port,
"tls_enabled": bool(node.tls_enabled), "tls_enabled": bool(node.tls_enabled),
@@ -96,6 +98,7 @@ def _build_hmac_proxy_url(ip: str, port: int, node_id: str, *, tls_enabled: bool
当 tls_enabled=True 时使用 https:// scheme。 当 tls_enabled=True 时使用 https:// scheme。
""" """
if not config.proxy_hmac_key: if not config.proxy_hmac_key:
logger.error("PROXY_HMAC_KEY 未配置,无法使用 ProxyNode 代理 (node_id={})", node_id)
raise ProxyNodeUnavailableError( raise ProxyNodeUnavailableError(
"PROXY_HMAC_KEY 未配置,无法使用 ProxyNode 代理", node_id=node_id "PROXY_HMAC_KEY 未配置,无法使用 ProxyNode 代理", node_id=node_id
) )
@@ -262,6 +265,7 @@ def build_proxy_url(proxy_config: dict[str, Any]) -> str | None:
node_id = node_id.strip() node_id = node_id.strip()
node_info = _get_proxy_node_info(node_id) node_info = _get_proxy_node_info(node_id)
if not node_info: if not node_info:
logger.warning("代理节点不可用(离线或不存在): node_id={}", node_id)
raise ProxyNodeUnavailableError(f"代理节点 {node_id} 不可用", node_id=node_id) raise ProxyNodeUnavailableError(f"代理节点 {node_id} 不可用", node_id=node_id)
# 手动节点:直接使用存储的代理 URL含认证信息 # 手动节点:直接使用存储的代理 URL含认证信息
@@ -325,6 +329,61 @@ def build_proxy_url(proxy_config: dict[str, Any]) -> str | None:
return proxy_url return proxy_url
def resolve_proxy_info(proxy_config: dict[str, Any] | None) -> dict[str, Any] | None:
"""
解析代理配置的摘要信息(用于日志和 usage 记录)
不构建实际的代理 URL仅返回可读的代理标识信息。
Returns:
{"node_id": "xxx", "node_name": "proxy-01", "source": "provider"} 或
{"url": "socks5://host:port", "source": "provider"} 或
{"node_id": "xxx", "node_name": "...", "source": "system"} 或
None (直连)
"""
source = "provider"
effective_config = proxy_config
# 无 provider 级代理时,尝试系统默认代理
if not effective_config or not effective_config.get("enabled", True):
effective_config = get_system_proxy_config()
source = "system"
if not effective_config or not effective_config.get("enabled", True):
return None
# ProxyNode 模式
node_id = effective_config.get("node_id")
if isinstance(node_id, str) and node_id.strip():
node_id = node_id.strip()
node_info = _get_proxy_node_info(node_id)
node_name = node_info.get("name", "unknown") if node_info else "offline"
return {"node_id": node_id, "node_name": node_name, "source": source}
# 旧格式 URL 模式
proxy_url = effective_config.get("url")
if proxy_url:
# 脱敏:只保留 scheme + host + port
try:
parsed = urlparse(proxy_url)
host_part = parsed.hostname or "unknown"
if parsed.port:
host_part = f"{host_part}:{parsed.port}"
safe_url = f"{parsed.scheme}://{host_part}"
except Exception:
safe_url = "unknown"
return {"url": safe_url, "source": source}
return None
def get_proxy_label(proxy_info: dict[str, Any] | None) -> str:
"""从 proxy_info 中提取简短的代理标签(用于日志)"""
if not proxy_info:
return "direct"
return proxy_info.get("node_name") or proxy_info.get("url") or "unknown"
def _make_proxy_param(proxy_url: str | None) -> str | httpx.Proxy | None: def _make_proxy_param(proxy_url: str | None) -> str | httpx.Proxy | None:
""" """
根据代理 URL 返回 httpx 可接受的 proxy 参数。 根据代理 URL 返回 httpx 可接受的 proxy 参数。

View File

@@ -746,9 +746,15 @@ class TaskService:
return "break" return "break"
if isinstance(cause, ProxyNodeUnavailableError): if isinstance(cause, ProxyNodeUnavailableError):
# ProxyNode 不可用属于配置明确指定但不可达/不可用的情况, # ProxyNode 不可用属于"配置明确指定但不可达/不可用"的情况,
# 在当前候选上重试通常没有意义,直接切换到下一个候选更合理。 # 在当前候选上重试通常没有意义,直接切换到下一个候选更合理。
logger.warning(" [{}] 代理节点不可用,切换候选: {}", request_id, str(cause)) node_id = cause.details.get("proxy_node_id") if cause.details else None
logger.warning(
" [{}] 代理节点不可用 (node_id={}),切换候选: {}",
request_id,
node_id or "unknown",
str(cause),
)
RequestCandidateService.mark_candidate_failed( RequestCandidateService.mark_candidate_failed(
db=self.db, db=self.db,
candidate_id=candidate_record_id, candidate_id=candidate_record_id,