diff --git a/ai/huawei_gateway.py b/ai/huawei_gateway.py index adb5bfa..8695925 100755 --- a/ai/huawei_gateway.py +++ b/ai/huawei_gateway.py @@ -412,58 +412,101 @@ def auto_compact(raw_body, real_token, orig_request): except Exception: final_summary = "\n".join(summaries) - logger.info(f"auto_compact: 完成, {body_size//1024}KB → {len(final_summary)}字 ({num_batches}批)") + logger.info(f"auto_compact: 压缩完成, {body_size//1024}KB → {len(final_summary)}字 ({num_batches}批)") - # 返回标准 OpenAI 格式(客户端无感知) - if stream: - chat_id = str(uuid.uuid4()) - created = int(time.time()) + # ============ 用压缩后的摘要+用户最新问题,重新请求模型 ============ + # 提取用户最后一条消息 + last_user_msg = None + for m in reversed(convo_msgs): + if m.get('role') == 'user': + last_user_msg = m + break - def compact_sse(): - 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" + # 构建压缩后的请求:system + 摘要(作为assistant上下文) + 用户最新问题 + compact_messages = system_msgs + [ + {"role": "assistant", "content": f"[上下文压缩摘要]\n{final_summary}"} + ] + if last_user_msg: + compact_messages.append(last_user_msg) - 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" + compact_payload = { + "model": model, + "messages": compact_messages, + "max_tokens": max_tokens, + "stream": stream, + } + # 保留原始请求中的其他参数 + for k in ('temperature', 'top_p', 'presence_penalty', 'frequency_penalty'): + if k in payload: + compact_payload[k] = payload[k] - 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" + compact_body = _json.dumps(compact_payload, ensure_ascii=False).encode('utf-8') + compact_headers = { + 'Authorization': f'Bearer {real_token}', + 'Host': TARGET_HOST, + 'Content-Type': 'application/json', + } + target_url = f'https://{TARGET_HOST}/v2/chat/completions' - return Response(compact_sse(), status=200, headers={ - 'Content-Type': 'text/event-stream;charset=UTF-8', - 'Cache-Control': 'no-cache', - 'X-Auto-Compact': f'batches={num_batches},original_kb={body_size//1024}', - }) - else: - return { - "id": f"chatcmpl-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) - } - } + logger.info(f"auto_compact: 用压缩上下文重新请求模型 ({len(compact_body)//1024}KB, stream={stream})") + + try: + resp = http_session.request( + method='POST', url=target_url, headers=compact_headers, + data=compact_body, allow_redirects=False, + timeout=UPSTREAM_TIMEOUT_MAX, stream=True + ) + + # 401 重试 + if resp.status_code == 401: + resp.close() + cache.blacklist_current() + new_token = find_token_in_memory() + if new_token: + compact_headers['Authorization'] = f'Bearer {new_token}' + resp = http_session.request( + method='POST', url=target_url, headers=compact_headers, + data=compact_body, allow_redirects=False, + timeout=UPSTREAM_TIMEOUT_MAX, stream=True + ) + + if resp.status_code != 200: + try: + err_body = resp.content[:500] + logger.error(f"auto_compact: 重新请求模型失败: HTTP {resp.status_code} - {err_body.decode('utf-8', errors='replace')}") + except: + logger.error(f"auto_compact: 重新请求模型失败: HTTP {resp.status_code}") + resp.close() + # 降级:返回摘要 + return {"error": {"message": f"压缩后重新请求失败(HTTP {resp.status_code}),上下文摘要: {final_summary[:500]}", "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] + resp_headers.append(('X-Auto-Compact', f'batches={num_batches},original_kb={body_size//1024}')) + + def compact_stream(): + try: + for chunk in resp.iter_content(chunk_size=16384): + if chunk: + yield chunk + finally: + resp.close() + + return Response(compact_stream(), status=200, headers=resp_headers, direct_passthrough=True) + else: + # 非流式 + content = resp.content + resp.close() + # 在响应头中标记经过了压缩 + return Response(content, status=200, + headers={'Content-Type': 'application/json', 'X-Auto-Compact': f'batches={num_batches},original_kb={body_size//1024}'}) + + except Exception as e: + logger.error(f"auto_compact: 重新请求模型异常: {e}") + return {"error": {"message": f"压缩后请求异常: {str(e)}", "type": "server_error"}}, 500 # ================= 全局请求日志(捕获所有请求,包括404) =================