| """Convert public agentic trajectory datasets to OpenAI-messages SFT samples |
| matching the pi harness prompt surface (system prompt + read/bash/edit/write tools). |
| |
| Sources: |
| 1. MaxDevv/real-pi-coding-agent-traces-sessions (native pi tool traces) |
| 2. R2E-Gym/R2EGym-SFT-Trajectories (XML <function=...> format) |
| 3. TIGER-Lab/SWE-Next-SFT-Trajectories (XML, has tool role) |
| 4. SWE-Gym/OpenHands-SFT-Trajectories (XML format) |
| 5. nvidia/Nemotron-Terminal-Corpus skill_based_* (Terminus-2 JSON format) |
| |
| Output: HF dataset with columns {messages, tools, source} saved to disk. |
| Only raw public hub datasets are used. Length-filtered with the Qwen3.5 tokenizer. |
| """ |
|
|
| import argparse |
| import glob |
| import json |
| import os |
| import re |
| import sys |
|
|
| sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) |
| from pi_prompt import TOOLS, make_system_prompt |
|
|
| from datasets import Dataset, load_dataset |
|
|
| CORE_TOOLS = {"read", "bash", "edit", "write"} |
| MAX_TOOL_RESPONSE_CHARS = 30000 |
|
|
|
|
| |
| |
| |
|
|
| def _trunc(s: str, n: int = MAX_TOOL_RESPONSE_CHARS) -> str: |
| if len(s) <= n: |
| return s |
| keep = n // 2 |
| return s[:keep] + f"\n[... {len(s) - n} chars truncated ...]\n" + s[-keep:] |
|
|
|
|
| def _sample(messages, source, cwd="/testbed", date="2026-08-15"): |
| return { |
| "messages": [{"role": "system", "content": make_system_prompt(date=date, cwd=cwd)}] + messages, |
| "tools": json.dumps(TOOLS), |
| "source": source, |
| } |
|
|
|
|
| |
| |
| |
|
|
| def convert_pi_traces(hf_home: str, max_sessions: int | None = None): |
| snap = glob.glob( |
| os.path.join(hf_home, "hub", "datasets--MaxDevv--real-pi-coding-agent-traces-sessions", "snapshots", "*") |
| )[0] |
| files = sorted(glob.glob(os.path.join(snap, "*.jsonl"))) |
| if max_sessions: |
| files = files[:max_sessions] |
|
|
| samples = [] |
| stats = {"sessions": 0, "episodes": 0, "drop_tool": 0, "drop_empty": 0, |
| "drop_nostop": 0, "drop_bashexec": 0, "drop_role": 0} |
|
|
| for path in files: |
| stats["sessions"] += 1 |
| events = [] |
| with open(path) as f: |
| for line in f: |
| try: |
| events.append(json.loads(line)) |
| except json.JSONDecodeError: |
| continue |
| msgs = [e["message"] for e in events if e.get("type") == "message" and "message" in e] |
| |
| episodes = [] |
| cur = [] |
| for m in msgs: |
| if m.get("role") == "user" and cur: |
| episodes.append(cur) |
| cur = [] |
| cur.append(m) |
| if cur: |
| episodes.append(cur) |
|
|
| for ep in episodes: |
| out = [] |
| ok = True |
| has_tool = False |
| final_stop = False |
| for m in ep: |
| role = m.get("role") |
| if role == "user": |
| text = "".join(b.get("text", "") for b in m.get("content", []) if isinstance(b, dict)) |
| if not text.strip(): |
| ok = False |
| break |
| out.append({"role": "user", "content": text}) |
| elif role == "assistant": |
| content_parts, think_parts, tool_calls = [], [], [] |
| stop = m.get("stopReason") |
| for b in m.get("content", []): |
| if not isinstance(b, dict): |
| continue |
| bt = b.get("type") |
| if bt == "text": |
| content_parts.append(b.get("text", "")) |
| elif bt == "thinking": |
| think_parts.append(b.get("thinking", b.get("text", ""))) |
| elif bt == "toolCall": |
| name = b.get("name") |
| if name not in CORE_TOOLS: |
| ok = False |
| break |
| tool_calls.append({ |
| "id": b.get("id", f"call_{len(tool_calls)}"), |
| "type": "function", |
| "function": {"name": name, "arguments": json.dumps(b.get("arguments", {}))}, |
| }) |
| has_tool = True |
| if not ok: |
| break |
| msg = {"role": "assistant", "content": "".join(content_parts)} |
| if think_parts: |
| msg["reasoning_content"] = "".join(think_parts) |
| if tool_calls: |
| msg["tool_calls"] = tool_calls |
| if not msg["content"] and not tool_calls and not think_parts: |
| continue |
| out.append(msg) |
| if stop == "stop": |
| final_stop = True |
| elif role == "toolResult": |
| if m.get("toolName") not in CORE_TOOLS: |
| ok = False |
| break |
| text = "".join(b.get("text", "") for b in m.get("content", []) if isinstance(b, dict)) |
| out.append({ |
| "role": "tool", |
| "tool_call_id": m.get("toolCallId", "call_0"), |
| "content": _trunc(text), |
| }) |
| elif role == "bashExecution": |
| ok = False |
| stats["drop_bashexec"] += 1 |
| break |
| else: |
| stats["drop_role"] += 1 |
| if not ok: |
| stats["drop_tool"] += 1 |
| continue |
| if not has_tool: |
| stats["drop_empty"] += 1 |
| continue |
| if not final_stop or not out or out[-1]["role"] != "assistant": |
| stats["drop_nostop"] += 1 |
| continue |
| stats["episodes"] += 1 |
| samples.append(_sample(out, "pi_traces", cwd="/app")) |
| return samples, stats |
|
|
|
|
| |
| |
| |
|
|
| FUNC_RE = re.compile(r"<function=([a-zA-Z_]+)>\s*(.*?)\s*</function>", re.DOTALL) |
| PARAM_RE = re.compile(r"<parameter=([a-zA-Z_]+)>\s*(.*?)\s*</parameter>", re.DOTALL) |
|
|
|
|
| def parse_xml_call(text: str): |
| """Return (pre_text, func_name, params) or None.""" |
| m = FUNC_RE.search(text) |
| if not m: |
| return None |
| pre = text[: m.start()].strip() |
| fname = m.group(1) |
| params = {p.group(1): p.group(2) for p in PARAM_RE.finditer(m.group(2))} |
| return pre, fname, params |
|
|
|
|
| def map_xml_tool(fname: str, params: dict): |
| """Map XML tool to (name, args) for pi tools, or 'finish', or None to drop episode.""" |
| if fname in ("file_editor", "str_replace_editor"): |
| cmd = params.get("command", "") |
| path = params.get("path", "") |
| if cmd == "view": |
| return ("read", {"path": path}) |
| if cmd == "create": |
| return ("write", {"path": path, "content": params.get("file_text", "")}) |
| if cmd == "str_replace": |
| return ("edit", {"path": path, "edits": [{"oldText": params.get("old_str", ""), "newText": params.get("new_str", "")}]}) |
| return None |
| if fname == "execute_bash": |
| return ("bash", {"command": params.get("cmd", params.get("command", ""))}) |
| if fname == "search": |
| term = params.get("search_term", "").replace("'", "'\\''") |
| path = params.get("path", ".") |
| return ("bash", {"command": f"grep -rn '{term}' {path} | head -50"}) |
| if fname == "finish": |
| return ("finish", {}) |
| return None |
|
|
|
|
| RESULT_PREFIX_RE = re.compile( |
| r"^(Execution output of \[[^\]]+\]|EXECUTION RESULT of \[[^\]]+\]|Exit code: \d+)\s*:?\s*\n?", |
| re.IGNORECASE, |
| ) |
|
|
|
|
| def strip_result_prefix(text: str) -> str: |
| prev = None |
| while prev != text: |
| prev = text |
| text = RESULT_PREFIX_RE.sub("", text) |
| return text.strip() |
|
|
|
|
| def convert_xml_traces(dataset_name: str, split: str, source: str, max_samples: int | None = None, |
| config: str | None = None): |
| ds = load_dataset(dataset_name, config, split=split) if config else load_dataset(dataset_name, split=split) |
| samples = [] |
| stats = {"in": len(ds), "kept": 0, "drop_parse": 0, "drop_tool": 0, "drop_flow": 0} |
| call_id = 0 |
|
|
| for ex in ds: |
| raw = ex["messages"] |
| |
| if raw and raw[0]["role"] == "system": |
| raw = raw[1:] |
| out = [] |
| ok = True |
| i = 0 |
| pending_call = None |
| while i < len(raw): |
| m = raw[i] |
| role = m["role"] |
| if role == "user" or role == "tool": |
| content = m["content"] |
| if pending_call is not None: |
| |
| text = strip_result_prefix(content) |
| out.append({"role": "tool", "tool_call_id": pending_call, "content": _trunc(text)}) |
| pending_call = None |
| else: |
| out.append({"role": "user", "content": content}) |
| i += 1 |
| elif role == "assistant": |
| parsed = parse_xml_call(m["content"]) |
| if parsed is None: |
| |
| if m["content"].strip(): |
| out.append({"role": "assistant", "content": m["content"].strip()}) |
| i += 1 |
| continue |
| pre, fname, params = parsed |
| mapped = map_xml_tool(fname, params) |
| if mapped is None: |
| ok = False |
| stats["drop_tool"] += 1 |
| break |
| name, args = mapped |
| if name == "finish": |
| text = pre or "The task is complete." |
| out.append({"role": "assistant", "content": text}) |
| i += 1 |
| continue |
| call_id += 1 |
| cid = f"call_{call_id}" |
| msg = {"role": "assistant", "content": pre, |
| "tool_calls": [{"id": cid, "type": "function", |
| "function": {"name": name, "arguments": json.dumps(args)}}]} |
| out.append(msg) |
| pending_call = cid |
| i += 1 |
| else: |
| i += 1 |
| if not ok: |
| continue |
| |
| if pending_call is not None or not out or out[-1]["role"] != "assistant": |
| stats["drop_flow"] += 1 |
| continue |
| if not any("tool_calls" in mm for mm in out): |
| stats["drop_flow"] += 1 |
| continue |
| stats["kept"] += 1 |
| samples.append(_sample(out, source, cwd="/testbed")) |
| if max_samples and len(samples) >= max_samples: |
| break |
| return samples, stats |
|
|
|
|
| |
| |
| |
|
|
| def convert_terminus(config: str, source: str, max_samples: int, seed: int = 0): |
| ds = load_dataset("nvidia/Nemotron-Terminal-Corpus", config, split="train") |
| if max_samples and len(ds) > max_samples: |
| ds = ds.shuffle(seed=seed).select(range(max_samples)) |
| samples = [] |
| stats = {"in": len(ds), "kept": 0, "drop_task": 0, "drop_parse": 0, "drop_flow": 0} |
| call_id = 0 |
|
|
| for ex in ds: |
| conv = ex["conversations"] |
| if not conv or conv[0]["role"] != "user": |
| stats["drop_flow"] += 1 |
| continue |
| first = conv[0]["content"] |
| tm = re.search(r"\n\nTask Description:\n", first) |
| if not tm: |
| stats["drop_task"] += 1 |
| continue |
| task = first[tm.end():] |
| |
| task = re.split(r"\nCurrent terminal state:|\nCurrent Terminal Screen:", task)[0].strip() |
| if not task: |
| stats["drop_task"] += 1 |
| continue |
|
|
| out = [{"role": "user", "content": task}] |
| ok = True |
| i = 1 |
| pending_call = None |
| while i < len(conv): |
| m = conv[i] |
| if m["role"] == "assistant": |
| content = m["content"] |
| think = "" |
| tm2 = re.match(r"\s*<think>(.*?)</think>\s*", content, re.DOTALL) |
| rest = content |
| if tm2: |
| think = tm2.group(1).strip() |
| rest = content[tm2.end():] |
| |
| jstart = rest.find("{") |
| if jstart == -1: |
| |
| if rest.strip(): |
| out.append({"role": "assistant", "content": rest.strip()}) |
| i += 1 |
| continue |
| try: |
| obj, _ = json.JSONDecoder().raw_decode(rest[jstart:]) |
| except json.JSONDecodeError: |
| ok = False |
| stats["drop_parse"] += 1 |
| break |
| analysis = (obj.get("analysis") or "").strip() |
| plan = (obj.get("plan") or "").strip() |
| commands = obj.get("commands") or [] |
| task_complete = obj.get("task_complete", False) |
| text = "\n\n".join(p for p in (analysis, f"Plan: {plan}" if plan else "") if p) |
| msg = {"role": "assistant", "content": text} |
| if think: |
| msg["reasoning_content"] = think |
| if commands: |
| script = "".join(c.get("keystrokes", "") for c in commands if isinstance(c, dict)) |
| if script.strip(): |
| call_id += 1 |
| cid = f"call_{call_id}" |
| msg["tool_calls"] = [{"id": cid, "type": "function", |
| "function": {"name": "bash", "arguments": json.dumps({"command": script})}}] |
| pending_call = cid |
| out.append(msg) |
| i += 1 |
| if task_complete: |
| break |
| elif m["role"] == "user": |
| content = m["content"] |
| content = re.sub(r"^New Terminal Output:\n?", "", content) |
| if pending_call is not None: |
| out.append({"role": "tool", "tool_call_id": pending_call, "content": _trunc(content)}) |
| pending_call = None |
| else: |
| out.append({"role": "user", "content": content}) |
| i += 1 |
| else: |
| i += 1 |
| if not ok: |
| continue |
| if pending_call is not None or not out or out[-1]["role"] != "assistant": |
| stats["drop_flow"] += 1 |
| continue |
| if not any("tool_calls" in mm for mm in out): |
| stats["drop_flow"] += 1 |
| continue |
| stats["kept"] += 1 |
| samples.append(_sample(out, source, cwd="/app")) |
| return samples, stats |
|
|
|
|
| |
| |
| |
|
|
| def _args_to_dict(messages): |
| out = [] |
| for m in messages: |
| m = dict(m) |
| if m.get("tool_calls"): |
| tcs = [] |
| for tc in m["tool_calls"]: |
| tc = dict(tc) |
| fn = dict(tc["function"]) |
| if isinstance(fn["arguments"], str): |
| fn["arguments"] = json.loads(fn["arguments"]) |
| tc["function"] = fn |
| tcs.append(tc) |
| m["tool_calls"] = tcs |
| out.append(m) |
| return out |
|
|
|
|
| def token_len_filter(samples, tokenizer, max_tokens: int): |
| kept, dropped = [], 0 |
| for s in samples: |
| try: |
| r = tokenizer.apply_chat_template( |
| _args_to_dict(s["messages"]), tools=json.loads(s["tools"]), add_generation_prompt=False |
| ) |
| n = len(r["input_ids"]) |
| except Exception: |
| dropped += 1 |
| continue |
| if n <= max_tokens: |
| s["n_tokens"] = n |
| kept.append(s) |
| else: |
| dropped += 1 |
| return kept, dropped |
|
|
|
|
| def main(): |
| ap = argparse.ArgumentParser() |
| ap.add_argument("--out", required=True) |
| ap.add_argument("--max-tokens", type=int, default=15000) |
| ap.add_argument("--max-pi-sessions", type=int, default=None) |
| ap.add_argument("--max-terminus", type=int, default=8000) |
| ap.add_argument("--max-r2e", type=int, default=None) |
| ap.add_argument("--max-swenext", type=int, default=None) |
| ap.add_argument("--skip-terminus", action="store_true") |
| ap.add_argument("--skip-pi", action="store_true") |
| ap.add_argument("--skip-xml", action="store_true") |
| args = ap.parse_args() |
|
|
| hf_home = os.environ["HF_HOME"] |
| all_samples = [] |
|
|
| if not args.skip_pi: |
| s, st = convert_pi_traces(hf_home, args.max_pi_sessions) |
| print(f"[pi_traces] {st}", flush=True) |
| all_samples += s |
|
|
| if not args.skip_xml: |
| s, st = convert_xml_traces("R2E-Gym/R2EGym-SFT-Trajectories", "train", "r2egym", args.max_r2e) |
| print(f"[r2egym] {st}", flush=True) |
| all_samples += s |
| s, st = convert_xml_traces("TIGER-Lab/SWE-Next-SFT-Trajectories", "train", "swenext", args.max_swenext) |
| print(f"[swenext] {st}", flush=True) |
| all_samples += s |
| s, st = convert_xml_traces("SWE-Gym/OpenHands-SFT-Trajectories", "train.success.oss", "openhands") |
| print(f"[openhands] {st}", flush=True) |
| all_samples += s |
|
|
| if not args.skip_terminus: |
| for cfg in ("skill_based_easy", "skill_based_medium", "skill_based_mixed"): |
| s, st = convert_terminus(cfg, f"terminus_{cfg}", args.max_terminus // 3) |
| print(f"[{cfg}] {st}", flush=True) |
| all_samples += s |
|
|
| print(f"total before length filter: {len(all_samples)}", flush=True) |
|
|
| from transformers import AutoTokenizer |
| tok = AutoTokenizer.from_pretrained("Qwen/Qwen3.5-9B-Base") |
| all_samples, dropped = token_len_filter(all_samples, tok, args.max_tokens) |
| print(f"length filter: kept {len(all_samples)}, dropped {dropped}", flush=True) |
|
|
| import random |
| random.seed(0) |
| random.shuffle(all_samples) |
| ds = Dataset.from_list(all_samples) |
| ds.save_to_disk(args.out) |
| print(f"saved {len(ds)} samples to {args.out}", flush=True) |
| from collections import Counter |
| print(Counter(ds["source"]), flush=True) |
| lens = ds["n_tokens"] |
| import numpy as np |
| print(f"tokens: mean {np.mean(lens):.0f} p50 {np.percentile(lens,50):.0f} p90 {np.percentile(lens,90):.0f} max {max(lens)}", flush=True) |
|
|
|
|
| if __name__ == "__main__": |
| main() |
|
|