Files
2026-09-04 14:58:42 +08:00

235 lines
8.6 KiB
Python

from __future__ import annotations
import json
import re
from pathlib import Path
from typing import Any, Iterable
from ..models import RolloutTrace
_TOOL_CALL_RE = re.compile(r"^Tool call (?P<name>[^:\n]+):[ \t]*", re.MULTILINE)
_TOOL_RESULT_RE = re.compile(
r"^Tool result(?: \((?P<name>[^;\n)]+)(?:;[ \t]*(?P<status>[^)\n]+))?\))?:?[ \t]*",
re.MULTILINE,
)
def _content(value: Any) -> str:
if value is None:
return ""
if isinstance(value, str):
return value
return json.dumps(value, ensure_ascii=False)
def opencode_events_to_state(
lines: Iterable[str], probe: str
) -> tuple[list[dict[str, str]], bool, int]:
state: list[dict[str, str]] = [{"role": "user", "content": probe}]
invoked = False
parsed = 0
for line in lines:
line = line.strip()
if not line:
continue
try:
event = json.loads(line)
except json.JSONDecodeError:
continue
if not isinstance(event, dict):
continue
parsed += 1
kind = str(event.get("type", event.get("event", ""))).lower()
part = event.get("part") if isinstance(event.get("part"), dict) else event
tool = part.get("tool") or part.get("name") or event.get("tool") or event.get("name")
title = part.get("title") or event.get("title") or ""
state_data = part.get("state") if isinstance(part.get("state"), dict) else {}
haystack = " ".join([str(kind), str(tool or ""), str(title), _content(part)])
if str(tool or "").lower() == "skill" or "<skill_content" in haystack.lower():
invoked = True
if tool or "tool" in kind:
arguments = state_data.get("input", part.get("input", part.get("arguments", {})))
output = state_data.get("output", part.get("output", part.get("result", "")))
state.append({"role": "assistant", "content": f"Tool call {tool or title}: {_content(arguments)}"})
if output not in (None, ""):
state.append({"role": "user", "content": f"Tool result: {_content(output)}"})
continue
text = part.get("text", part.get("content", event.get("message", "")))
if text not in (None, ""):
role = str(event.get("role", part.get("role", "assistant")))
if role not in {"assistant", "user", "system"}:
role = "assistant"
state.append({"role": role, "content": _content(text)})
return state, invoked, parsed
def read_event_file(path: Path, probe: str) -> tuple[list[dict[str, str]], bool, int]:
with path.open(encoding="utf-8", errors="replace") as handle:
return opencode_events_to_state(handle, probe)
def _split_tool_results(content: str) -> tuple[str, list[tuple[str, str, str]]]:
matches = list(_TOOL_RESULT_RE.finditer(content))
if not matches:
return content, []
prefix = content[:matches[0].start()].strip()
results = []
for index, match in enumerate(matches):
end = matches[index + 1].start() if index + 1 < len(matches) else len(content)
results.append((
(match.group("name") or "unknown").strip(),
(match.group("status") or "unknown").strip(),
content[match.end():end].strip(),
))
return prefix, results
def _excerpt(content: str, limit: int) -> str:
content = content.strip()
if len(content) <= limit:
return content
marker = "\n[... content omitted ...]\n"
if limit <= len(marker) + 2:
return content[:limit]
omitted = len(content) - (limit - len(marker))
while True:
marker = f"\n[... {omitted} chars omitted ...]\n"
available = limit - len(marker)
updated = len(content) - available
if updated == omitted:
break
omitted = updated
head = (available + 1) // 2
tail = available // 2
return content[:head] + marker + content[-tail:]
def _termination(trace: RolloutTrace) -> str:
value = trace.metadata.get("termination")
if value not in (None, ""):
return str(value)
if trace.timed_out is True:
return "timeout"
if trace.exit_code == 0:
return "completed"
if trace.exit_code is not None:
return "error"
return "unknown"
def _runtime_facts(trace: RolloutTrace) -> str:
metadata = trace.metadata
facts = {
"trace_id": trace.trace_id,
"termination": _termination(trace),
"timed_out": trace.timed_out,
"exit_code": trace.exit_code,
"agent_execution_seconds": metadata.get("agent_execution_seconds"),
"tool_calls": metadata.get("tool_calls"),
"skill_invoked": trace.skill_invoked,
"timeout_reason": metadata.get("timeout_reason"),
"error_category": metadata.get("error_category"),
"partial_trajectory": metadata.get("partial_trajectory"),
}
return "RUNTIME_FACTS " + json.dumps(facts, ensure_ascii=False, separators=(",", ":"))
def _allocate_excerpt_budget(caps: list[int], available: int) -> list[int]:
allocations = [0] * len(caps)
active = [index for index, cap in enumerate(caps) if cap > 0]
while active and available > 0:
share = max(1, available // len(active))
progressed = False
for index in active.copy():
amount = min(share, caps[index] - allocations[index], available)
allocations[index] += amount
available -= amount
progressed = progressed or amount > 0
if allocations[index] >= caps[index]:
active.remove(index)
if available == 0:
break
if not progressed:
break
return allocations
def compact_trace(trace: RolloutTrace, total: int = 30000) -> str:
"""Render runtime facts and a complete action ledger for one Map call."""
entries: list[dict[str, Any]] = []
for index, message in enumerate(trace.state):
event_id = f"E{index:03d}"
role = str(message.get("role", "unknown"))
content = _content(message.get("content", ""))
tool_call = _TOOL_CALL_RE.match(content)
if tool_call:
entries.append({
"id": f"{event_id}.T01",
"kind": "tool_call",
"role": role,
"name": tool_call.group("name").strip(),
"status": "unknown",
"content": content[tool_call.end():].strip(),
})
continue
prefix, tool_results = _split_tool_results(content) if index > 0 else (content, [])
if prefix:
entries.append({
"id": event_id,
"kind": "message",
"role": role,
"content": prefix,
})
for tool_index, (name, status, result) in enumerate(tool_results, 1):
entries.append({
"id": f"{event_id}.T{tool_index:02d}",
"kind": "tool_result",
"role": role,
"name": name,
"status": status,
"content": result,
})
message_entries = [entry for entry in entries if entry["kind"] == "message"]
first_user = next((entry for entry in message_entries if entry["role"] == "user"), None)
final_assistant = next(
(entry for entry in reversed(message_entries) if entry["role"] == "assistant"),
None,
)
skeletons = []
caps = []
for entry in entries:
if entry["kind"] == "message":
labels = ["message", f"role={entry['role']}"]
if entry is first_user:
labels.append("task")
if entry is final_assistant:
labels.append("final")
skeleton = f"{entry['id']} [" + " ".join(labels) + "]"
cap = 2500 if entry is first_user else 3000 if entry is final_assistant else 600
else:
name = str(entry["name"])[:80]
status = str(entry["status"])[:40]
skeleton = f"{entry['id']} [{entry['kind']} name={name} status={status}]"
cap = 600
skeletons.append(skeleton)
caps.append(min(cap, len(str(entry["content"]))))
facts = _runtime_facts(trace)
fixed_size = (
len(facts)
+ sum(len(skeleton) + 1 for skeleton in skeletons)
+ sum(1 for cap in caps if cap > 0)
)
allocations = _allocate_excerpt_budget(caps, max(0, total - fixed_size))
rendered = [facts]
for entry, skeleton, allocation in zip(entries, skeletons, allocations):
rendered.append(skeleton)
if allocation:
rendered.append(_excerpt(str(entry["content"]), allocation))
return "\n".join(rendered)