From 85ea33eb4e75fef070df5f17599348b0571439c7 Mon Sep 17 00:00:00 2001 From: chaos Date: Wed, 8 Jul 2026 14:06:21 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E8=A7=A3=E5=86=B3=20compact=20=E5=91=BD?= =?UTF-8?q?=E4=BB=A4=E6=96=AD=E8=BF=9E=E9=97=AE=E9=A2=98=20-=20=E8=AF=B7?= =?UTF-8?q?=E6=B1=82=E4=BD=93=E8=B6=85=E9=99=90413=E6=B8=85=E6=99=B0?= =?UTF-8?q?=E9=94=99=E8=AF=AF=20+=20/v2/compact=E5=88=86=E6=89=B9=E5=8E=8B?= =?UTF-8?q?=E7=BC=A9=E7=AB=AF=E7=82=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 根因: APIG请求体限制~1.25MB, 超限后JSON截断导致'model id缺失'400错误→客户端断连 修复: 1. 请求体超限校验: 返回413 + content_too_large清晰错误(而非上游的'model id缺失') 2. /v2/compact端点: 超长对话自动分批发送压缩,每批<800KB,最后合并总结 3. 动态超时: 根据请求体大小自动调整(60s→180s→300s) 4. 修复响应头转发: 剥离Content-Encoding避免客户端解压失败 5. 禁用flask-compress(与Waitress SSE兼容问题) --- ai/huawei_gateway.py | 316 +++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 308 insertions(+), 8 deletions(-) diff --git a/ai/huawei_gateway.py b/ai/huawei_gateway.py index 3cb7005..7dbef8c 100755 --- a/ai/huawei_gateway.py +++ b/ai/huawei_gateway.py @@ -14,6 +14,7 @@ import os import re import sys import time +import gzip import logging import threading import traceback @@ -42,6 +43,12 @@ MAX_MEM_SEGMENT = 200 * 1024 * 1024 # 单段最大扫描 200MB TOKEN_PATTERN = re.compile(b'Bearer ([A-Za-z0-9+/=_-]{100,})') TARGET_HOST = 'tokenhub.developer.huaweicloud.com' +# ================= 请求体限制配置 ================= +APIG_BODY_LIMIT = 1200 * 1024 # APIG 请求体限制 ~1.2MB (实测边界1260KB) +GZIP_THRESHOLD = 200 * 1024 # 超过 200KB 时启用 gzip 压缩转发 +UPSTREAM_TIMEOUT_MIN = 60 # 小请求超时 60s +UPSTREAM_TIMEOUT_MAX = 300 # 大请求超时 300s + # ================= 日志 ================= logging.basicConfig( level=logging.INFO, @@ -248,13 +255,12 @@ http_session.mount('http://', adapter) # ================= Flask 应用 ================= app = Flask(__name__) -# 启用响应压缩(gzip/brotli/zstd),对长文本响应压缩率可达 70-80% -if _compress_available: +# 响应压缩(flask-compress 与 Waitress SSE 存在兼容问题,暂时禁用) +# 如需启用,需切换到 Gunicorn 或确保客户端不发送 Accept-Encoding +if False and _compress_available: compress = Compress() compress.init_app(app) - # 设置压缩最小阈值(256字节以下不压缩) app.config['COMPRESS_MIN_SIZE'] = 256 - # 压缩 MIME 类型白名单 app.config['COMPRESS_MIMETYPES'] = [ 'text/html', 'text/plain', 'text/css', 'text/xml', 'application/json', 'application/javascript', @@ -291,6 +297,276 @@ def set_token(): return {"status": "ok", "token_fingerprint": cache._fingerprint(token)}, 200 +# ================= Compact 分批压缩端点 ================= +@app.route('/v2/compact', methods=['POST', 'OPTIONS']) +def compact_endpoint(): + """ + 分批上下文压缩端点: + - 对话历史超长时,自动分批发送给模型压缩 + - 每批独立总结,最后合并为完整压缩上下文 + - 兼容 OpenAI chat completions 请求格式 + """ + if request.method == 'OPTIONS': + return Response(status=200, headers={ + 'Access-Control-Allow-Origin': '*', + 'Access-Control-Allow-Methods': 'POST, OPTIONS', + 'Access-Control-Allow-Headers': 'Content-Type, Authorization' + }) + + import json as _json + + real_token = find_token_in_memory() + if not real_token: + return {"error": {"message": "未找到华为云Token", "type": "server_error"}}, 500 + + try: + payload = request.get_json(force=True) + except Exception: + return {"error": {"message": "无效的JSON请求体", "type": "invalid_request_error"}}, 400 + + model = payload.get('model', 'glm-5.1') + messages = payload.get('messages', []) + max_tokens = payload.get('max_tokens', 500) + stream = payload.get('stream', False) + + if not messages: + return {"error": {"message": "messages 不能为空", "type": "invalid_request_error"}}, 400 + + # 计算实际请求体大小(使用原始请求体,而非重新序列化) + raw_body = request.get_data() + body_size = len(raw_body) + + # 如果请求体在限制内,直接转发给上游(无需分批) + if body_size <= APIG_BODY_LIMIT: + logger.info(f"compact: 请求体{body_size//1024}KB在限制内,直接转发") + target_url = f'https://{TARGET_HOST}/v2/chat/completions' + headers = { + 'Authorization': f'Bearer {real_token}', + 'Host': TARGET_HOST, + 'Content-Type': 'application/json', + } + # 保留客户端的其他头 + for k, v in request.headers: + kl = k.lower() + if kl not in ('host', 'content-length', 'connection', 'accept-encoding', + 'transfer-encoding', 'authorization', 'content-type'): + headers[k] = v + + try: + resp = http_session.request( + method='POST', url=target_url, headers=headers, + data=raw_body, allow_redirects=False, + timeout=UPSTREAM_TIMEOUT_MAX, stream=True + ) + except Exception as e: + logger.error(f"compact 转发失败: {e}") + return {"error": {"message": f"转发失败: {str(e)}", "type": "server_error"}}, 502 + + # 流式转发 + content_type = resp.headers.get('Content-Type', '') + if 'text/event-stream' in content_type or stream: + skip_h = {'transfer-encoding', 'content-encoding', 'content-length', + 'connection', 'keep-alive', 'upgrade'} + resp_headers = [(k, v) for k, v in resp.headers.items() if k.lower() not in skip_h] + + def sse_stream(): + try: + for chunk in resp.iter_content(chunk_size=16384): + if chunk: + yield chunk + finally: + resp.close() + return Response(sse_stream(), status=resp.status_code, + headers=resp_headers, direct_passthrough=True) + else: + content = resp.content + resp.close() + skip_h2 = {'transfer-encoding', 'content-encoding', 'content-length', + 'connection', 'keep-alive', 'upgrade'} + return Response(content, status=resp.status_code, + headers=[(k, v) for k, v in resp.headers.items() + if k.lower() not in skip_h2]) + + # ============ 分批压缩逻辑 ============ + logger.info(f"compact: 请求体{body_size//1024}KB超限,启动分批压缩") + + # 分离 system 消息和对话消息 + system_msgs = [m for m in messages if m.get('role') == 'system'] + convo_msgs = [m for m in messages if m.get('role') != 'system'] + + if not convo_msgs: + return {"error": {"message": "无对话内容可压缩", "type": "invalid_request_error"}}, 400 + + # 按批次分割对话:每批控制在安全大小内 + SAFE_BATCH_BYTES = 800 * 1024 # 每批 800KB(留余量给 JSON 开销) + batches = [] + current_batch = [] + current_size = 0 + + for msg in convo_msgs: + msg_size = len(_json.dumps(msg, ensure_ascii=False).encode('utf-8')) + if current_size + msg_size > SAFE_BATCH_BYTES and current_batch: + batches.append(current_batch) + current_batch = [msg] + current_size = msg_size + else: + current_batch.append(msg) + current_size += msg_size + + if current_batch: + batches.append(current_batch) + + logger.info(f"compact: 分为{len(batches)}批, 每批~{SAFE_BATCH_BYTES//1024}KB") + + # 逐批压缩 + target_url = f'https://{TARGET_HOST}/v2/chat/completions' + summaries = [] + + for i, batch in enumerate(batches): + # 构建压缩请求 + batch_messages = system_msgs + batch + [{ + "role": "user", + "content": "请对以上对话内容进行简洁压缩总结,保留所有关键信息、决策和结论,去除冗余和重复。用简洁的条目式格式输出。" + }] + + batch_payload = { + "model": model, + "messages": batch_messages, + "max_tokens": max_tokens, + "stream": False, + "temperature": 0.3 + } + + batch_body = _json.dumps(batch_payload, ensure_ascii=False).encode('utf-8') + + headers = { + 'Authorization': f'Bearer {real_token}', + 'Host': TARGET_HOST, + 'Content-Type': 'application/json', + } + + try: + resp = http_session.request( + method='POST', url=target_url, headers=headers, + data=batch_body, allow_redirects=False, + timeout=UPSTREAM_TIMEOUT_MAX, stream=False + ) + data = resp.json() + + if resp.status_code == 200 and 'choices' in data: + summary = data['choices'][0].get('message', {}).get('content', '') + summaries.append(summary) + logger.info(f"compact: 第{i+1}/{len(batches)}批压缩完成, {len(summary)}字") + else: + err_msg = data.get('error', {}).get('message', str(data)[:200]) + logger.warning(f"compact: 第{i+1}批失败: HTTP {resp.status_code} - {err_msg}") + # 失败的批次保留原始内容摘要 + fallback = f"[第{i+1}批压缩失败,原始{len(batch)}条消息]" + summaries.append(fallback) + except Exception as e: + logger.error(f"compact: 第{i+1}批异常: {e}") + summaries.append(f"[第{i+1}批压缩异常]") + + time.sleep(1) # 避免触发限流 + + # 合并所有批次的压缩结果 + if len(summaries) == 1: + final_summary = summaries[0] + else: + # 对多个摘要做最终合并压缩 + merge_messages = system_msgs + [{ + "role": "user", + "content": "以下是分批压缩的对话摘要,请合并为一个连贯的压缩总结,保留所有关键信息:\n\n" + + "\n\n---\n\n".join(f"第{i+1}批摘要:\n{s}" for i, s in enumerate(summaries)) + }] + + merge_payload = { + "model": model, + "messages": merge_messages, + "max_tokens": max_tokens, + "stream": False, + "temperature": 0.3 + } + + try: + resp = http_session.request( + method='POST', url=target_url, headers=headers, + data=_json.dumps(merge_payload, ensure_ascii=False).encode('utf-8'), + allow_redirects=False, timeout=UPSTREAM_TIMEOUT_MAX + ) + merge_data = resp.json() + if resp.status_code == 200 and 'choices' in merge_data: + final_summary = merge_data['choices'][0].get('message', {}).get('content', '') + else: + final_summary = "\n".join(summaries) + except Exception: + final_summary = "\n".join(summaries) + + logger.info(f"compact: 压缩完成, 最终{len(final_summary)}字 (原始{body_size//1024}KB)") + + # 返回 OpenAI 兼容格式 + if stream: + # 流式返回:将压缩结果包装为 SSE 事件 + import uuid + chat_id = str(uuid.uuid4()) + created = int(time.time()) + + def compact_sse(): + # 首个 chunk:role + first = { + "id": chat_id, "object": "chat.completion.chunk", "created": created, + "model": model, + "choices": [{"index": 0, "delta": {"role": "assistant", "content": ""}, "finish_reason": None}] + } + yield f"data: {_json.dumps(first, ensure_ascii=False)}\n\n" + + # 内容 chunk + content_chunk = { + "id": chat_id, "object": "chat.completion.chunk", "created": created, + "model": model, + "choices": [{"index": 0, "delta": {"content": final_summary}, "finish_reason": None}] + } + yield f"data: {_json.dumps(content_chunk, ensure_ascii=False)}\n\n" + + # 结束 chunk + done_chunk = { + "id": chat_id, "object": "chat.completion.chunk", "created": created, + "model": model, + "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}] + } + yield f"data: {_json.dumps(done_chunk, ensure_ascii=False)}\n\n" + yield "data: [DONE]\n\n" + + return Response(compact_sse(), status=200, headers={ + 'Content-Type': 'text/event-stream;charset=UTF-8', + 'Cache-Control': 'no-cache', + 'X-Compact-Batches': str(len(batches)), + 'X-Compact-Original-KB': str(body_size // 1024), + }, direct_passthrough=True) + else: + # 非流式返回 + return { + "id": f"compact-{int(time.time())}", + "object": "chat.completion", + "created": int(time.time()), + "model": model, + "choices": [{ + "index": 0, + "message": {"role": "assistant", "content": final_summary}, + "finish_reason": "stop" + }], + "usage": { + "prompt_tokens": body_size // 4, # 估算 + "completion_tokens": len(final_summary), + "total_tokens": body_size // 4 + len(final_summary) + }, + "compact_meta": { + "batches": len(batches), + "original_kb": body_size // 1024, + } + } + + @app.route('/v2/', methods=['POST', 'GET', 'OPTIONS', 'PUT', 'DELETE']) def proxy(subpath): if request.method == 'OPTIONS': @@ -317,16 +593,40 @@ def proxy(subpath): target_url = f'https://{TARGET_HOST}/v2/{subpath}' + # ============ 请求体处理:大小校验 + 动态超时 ============ + raw_body = request.get_data() + body_size = len(raw_body) + + # 动态超时:根据请求体大小自动调整 + if body_size > 800 * 1024: + upstream_timeout = UPSTREAM_TIMEOUT_MAX + elif body_size > 200 * 1024: + upstream_timeout = 180 + else: + upstream_timeout = UPSTREAM_TIMEOUT_MIN + + # 请求体超限校验:返回清晰错误而非上游的 "model id 缺失" + if body_size > APIG_BODY_LIMIT: + logger.warning(f"请求体超限: {body_size//1024}KB > {APIG_BODY_LIMIT//1024}KB (APIG限制)") + return { + "error": { + "message": f"请求体过大({body_size//1024}KB),超过API网关限制({APIG_BODY_LIMIT//1024}KB)。请减少对话历史长度,或使用 /v2/compact 端点进行分批压缩。", + "type": "invalid_request_error", + "code": "content_too_large", + "param": None + } + }, 413 + try: # 使用 stream=True 支持 SSE 流式转发 resp = http_session.request( method=request.method, url=target_url, headers=headers, - data=request.get_data(), + data=raw_body, cookies=request.cookies, allow_redirects=False, - timeout=60, + timeout=upstream_timeout, stream=True ) @@ -342,10 +642,10 @@ def proxy(subpath): method=request.method, url=target_url, headers=headers, - data=request.get_data(), + data=raw_body, cookies=request.cookies, allow_redirects=False, - timeout=60, + timeout=upstream_timeout, stream=True )