Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
103 changes: 40 additions & 63 deletions src/converter/anti_truncation.py
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,6 @@
现在请继续输出:"""



# ==================== 请求注入 ====================


Expand Down Expand Up @@ -126,8 +125,6 @@ def apply_anti_truncation(payload: Dict[str, Any]) -> Dict[str, Any]:
return modified_payload




# ==================== 响应提取 ====================


Expand Down Expand Up @@ -221,7 +218,6 @@ def _extract_content_from_json_str(args_str: str) -> Optional[str]:
return None



def build_text_chunk_from_synthetic(
original_data: Dict[str, Any],
synthetic_content: str,
Expand Down Expand Up @@ -288,8 +284,6 @@ def build_text_chunk_from_synthetic(
return modified_inner




# ==================== 流式处理器 ====================


Expand Down Expand Up @@ -354,7 +348,8 @@ async def process_stream(self) -> AsyncGenerator[bytes, None]:
# 每轮流内状态:初始化放在 try 之外,流中断的异常路径需要访问
found_synthetic = False
has_real_tool_calls = False
side_buffer = io.StringIO() # 暂存普通文本(防拼接)
emitted_text = False
side_buffer = io.StringIO() # 保存已透传正文,供异常中断后续传
last_finish_reason: Optional[str] = None

try:
Expand Down Expand Up @@ -444,14 +439,14 @@ async def process_stream(self) -> AsyncGenerator[bytes, None]:
synthetic_content, real_calls, chunk_has_synthetic = (
extract_synthetic_content_from_response(data)
)
has_real_tool_calls = has_real_tool_calls or bool(real_calls)

if chunk_has_synthetic:
found_synthetic = True
# 防拼接:丢弃之前暂存的普通文本
if side_buffer.getvalue():
if emitted_text:
print(
"Anti-truncation: Discarding side-buffered text "
"(content conflict with synthetic tool)",
"Anti-truncation: Plain text was streamed before "
"synthetic tool call (possible duplication)",
flush=True,
)
side_buffer.close()
Expand All @@ -471,7 +466,6 @@ async def process_stream(self) -> AsyncGenerator[bytes, None]:

elif real_calls:
# 真实工具调用,原样透传
has_real_tool_calls = True
yield line

else:
Expand All @@ -487,24 +481,25 @@ async def process_stream(self) -> AsyncGenerator[bytes, None]:
stripped, separators=(",", ":"), ensure_ascii=False
)
yield f"data: {json_str}\n\n".encode("utf-8")
elif not self._chunk_has_plain_text(data):
yield line
continue
else:
# 暂存到 side buffer,等待看是否有合成工具调用
text = self._extract_text_from_chunk(data)
if text:
side_buffer.write(text)
# 正文与思考实时透传;正文另存一份用于异常续传历史。
if self._chunk_has_plain_text(data):
emitted_text = True
side_buffer.write(self._extract_text_from_chunk(data))
chunk_finish = self._get_finish_reason(data)
if chunk_finish:
last_finish_reason = chunk_finish
# 暂时不透传,等流结束时决定
yield line
continue

else:
# 非 data: 开头的行,直接传递
yield line

# 流结束(break 或正常结束)
side_text = side_buffer.getvalue()
side_buffer.close()

if found_synthetic:
Expand All @@ -514,41 +509,17 @@ async def process_stream(self) -> AsyncGenerator[bytes, None]:
yield b"data: [DONE]\n\n"
return

# 模型自然结束(STOP)且未调用合成工具:视为完整回答,不再续传
# (续传会让模型重发一遍已有内容,浪费上游请求且可能产生重复文本)
if last_finish_reason == "STOP":
# 真实工具或正文已交给客户端,不能续传重放;保留上游 STOP 终止语义。
if has_real_tool_calls or emitted_text or last_finish_reason == "STOP":
print(
f"Anti-truncation: Stream ended with STOP without synthetic tool call "
f"(text length: {len(side_text)}), treating as complete",
"Anti-truncation: Real output or STOP received, treating turn as complete",
flush=True,
)
if side_text:
# 输出暂存文本(正常情况下模型守规矩时不会走到这里)
self._append_content(side_text)
fallback_chunk = self._build_fallback_text_chunk(side_text)
if fallback_chunk:
yield fallback_chunk
self._clear_content()
# 补发 finishReason 收尾 chunk(Gemini 客户端靠它判断流正常结束)
yield self._build_finish_reason_chunk()
yield b"data: [DONE]\n\n"
return

# 未收到合成工具调用
if side_text:
# 有普通文本作为 fallback,输出它
print(
f"Anti-truncation: No synthetic tool call, "
f"using side-buffered text as fallback (length: {len(side_text)})",
flush=True,
)
self._append_content(side_text)
# 构建一个包含 side buffer 文本的 chunk 输出
fallback_chunk = self._build_fallback_text_chunk(side_text)
if fallback_chunk:
yield fallback_chunk

# 触发续传
# 正常结束但没有正文或工具产出(如纯思考)才触发续传。
if self.current_attempt < self.max_attempts:
accumulated_text = self._get_collected_text()
total_length = len(accumulated_text)
Expand Down Expand Up @@ -580,20 +551,12 @@ async def process_stream(self) -> AsyncGenerator[bytes, None]:
)
self._append_content(interrupted_text)

if self.current_attempt >= self.max_attempts:
# 重试额度用尽:把已收到的内容(含残文)作为 fallback 输出,尽量不浪费
if has_real_tool_calls or self.current_attempt >= self.max_attempts:
# 工具已经发出时必须报错终止,续传可能导致客户端重复执行。
# 正文已实时发送;重试额度用尽时不能再次输出整个历史。
salvaged = self._get_collected_text()
self._clear_content()
if salvaged:
print(
f"Anti-truncation: Max attempts reached after error, "
f"yielding salvaged text (length: {len(salvaged)})",
flush=True,
)
fallback_chunk = self._build_fallback_text_chunk(salvaged)
if fallback_chunk:
yield fallback_chunk
else:
if has_real_tool_calls or not salvaged:
error_chunk = {
"error": {
"message": f"Anti-truncation failed: {str(e)}",
Expand Down Expand Up @@ -633,8 +596,8 @@ def _build_current_payload(self) -> Dict[str, Any]:
if accumulated_text:
new_contents.append({"role": "model", "parts": [{"text": accumulated_text}]})

# 预填充模式:直接用拼接内容作为末尾 model 预填充
if self.enable_prefill_mode:
# 没有正文可续写时使用明确的续传指令,避免空预填充原样重发。
if self.enable_prefill_mode and accumulated_text:
request_data["contents"] = new_contents
continuation_payload["request"] = request_data
return continuation_payload
Expand Down Expand Up @@ -669,10 +632,21 @@ def _extract_text_from_chunk(self, data: Dict[str, Any]) -> str:
for candidate in data.get("candidates", []):
content = candidate.get("content", {})
for part in content.get("parts", []):
if isinstance(part, dict) and "text" in part:
if isinstance(part, dict) and "text" in part and not part.get("thought", False):
text += part["text"]
return text

@staticmethod
def _chunk_has_plain_text(data: Dict[str, Any]) -> bool:
"""非空正文才算产出;思考与空文本收尾块不算。"""
if "response" in data:
data = data["response"]
return any(
isinstance(part, dict) and part.get("text") and not part.get("thought", False)
for candidate in data.get("candidates", [])
for part in candidate.get("content", {}).get("parts", [])
)

@staticmethod
def _has_finish_reason(data: Dict[str, Any]) -> bool:
"""判断 chunk 是否携带 finishReason(控制信号 chunk)。"""
Expand Down Expand Up @@ -822,6 +796,11 @@ async def _handle_non_streaming_response(self, response) -> bytes:
extract_synthetic_content_from_response(response_data)
)

if not found_synthetic and (
real_calls or self._chunk_has_plain_text(response_data)
):
return content.encode() if isinstance(content, str) else content

if found_synthetic or self.current_attempt >= self.max_attempts:
if found_synthetic:
# 替换响应中的合成工具调用为普通文本
Expand Down Expand Up @@ -856,5 +835,3 @@ async def _handle_non_streaming_response(self, response) -> bytes:
}
}
).encode()


Loading