"""Sub-goal annotation of real trajectories (Plan-and-Act style), with the SAME planner prompt the agent uses at run time. python finetune/annotate_subgoals.py finetune/out/m2w_cases.jsonl finetune/out/m2w_sub_cases.jsonl [max_tasks=0] [workers=12] Per task (grouped by task_id, steps ordered by history length): local Qwen (TEXT_MODEL_BASE_URL) segments the recorded steps, in order, into sub-goals written in the run-time planner's style (trajectory-conditioned, so plan and steps agree). Emits, per step, a copy of the case with goal = its sub-goal (same gold). With EMIT_DONE=1 it also emits, at every boundary, the next step's page with goal = the finished sub-goal and gold = DONE. Single-sub-goal tasks are skipped. """ import json, os, sys, threading from collections import defaultdict from concurrent.futures import ThreadPoolExecutor sys.path.insert(0, "/home/ckl/projects/S/jev-ultrafast") import importlib.util _spec = importlib.util.spec_from_file_location("jq", "/home/ckl/projects/S/jev-ultrafast/jev_ultrafast/questions.py") jq = importlib.util.module_from_spec(_spec); _spec.loader.exec_module(jq) import httpx URL = os.environ.get("TEXT_MODEL_BASE_URL", "http://127.0.0.1:30000/v1").rstrip("/") + "/chat/completions" MODEL = os.environ.get("TEXT_MODEL", "Qwen/Qwen3-8B-AWQ") CLIENT = httpx.Client(timeout=120, trust_env=False) # local server: never through a proxy # segment boundaries are right most of the time but not always; a wrong boundary turns into a wrong DONE label, so # sub-goal DONE cases come only from webgym (verified by page state) unless EMIT_DONE=1 EMIT_DONE = os.environ.get("EMIT_DONE") == "1" SEGMENT = """You annotate a recorded browser trajectory with the sub-goals a planner would have written for it. Split the recorded steps, in their recorded order, into consecutive segments; each segment is one sub-goal. Write each sub-goal as an instruction for the part of the task those steps accomplish, in the task's own words: name what is set and repeat the value (patterns: "Set the to ", "Search for and show the results", "Sort the results by ", "Open "). Steps that only open a menu or a picker belong to the sub-goal they serve; a sub-goal ends with the step that actually commits its value. A task done in one go is one segment. Return JSON {"segments": [{"subgoal": "...", "last_step": k}, ...]} where last_step is the 0-based index of the segment's last step; last_step values increase and the final one is the last step index.""" def chat(system, user, max_tokens=400): r = CLIENT.post(URL, json={"model": MODEL, "max_tokens": max_tokens, "temperature": 0, "response_format": {"type": "json_object"}, "chat_template_kwargs": {"enable_thinking": False}, "messages": [{"role": "system", "content": system}, {"role": "user", "content": user}]}) r.raise_for_status() return json.loads(r.json()["choices"][0]["message"]["content"]) def step_repr(c): t = f" '{c['history'][-1]['text']}'" if False else "" lab = (c.get("label") or "")[:80] return f"{c['gold_op']} {lab}".strip() def annotate(task_cases): goal = task_cases[0]["goal"] steps = [step_repr(c) for c in task_cases] # typed / selected values are only visible in the next case's history; attach them for context for i, c in enumerate(task_cases[:-1]): h = task_cases[i + 1]["history"][-1] if task_cases[i + 1]["history"] else None if h and h.get("text") and c["gold_op"] in ("TYPE_TEXT", "SELECT"): steps[i] += f" = '{h['text']}'" r = chat(SEGMENT, json.dumps({"task": goal, "steps": [f"{i}: {s}" for i, s in enumerate(steps)]}, ensure_ascii=False), max_tokens=600) segs = r.get("segments") or [] plan, a, start = [], [None] * len(task_cases), 0 for sg in segs: try: last = int(sg["last_step"]); text = str(sg["subgoal"]).strip() except Exception: return None if not text or last < start or last >= len(task_cases): return None for k in range(start, last + 1): a[k] = len(plan) plan.append(text); start = last + 1 if start != len(task_cases) or len(plan) < 2: return None return plan, a def main(): src, dst = sys.argv[1], sys.argv[2] max_tasks = int(sys.argv[3]) if len(sys.argv) > 3 else 0 workers = int(sys.argv[4]) if len(sys.argv) > 4 else 12 by = defaultdict(list) for l in open(src): c = json.loads(l); by[c.get("task_id") or c["goal"]].append(c) tasks = [sorted(v, key=lambda c: len(c["history"])) for v in by.values()] tasks = [t for t in tasks if len(t) >= 2] if max_tasks: tasks = tasks[:max_tasks] lock = threading.Lock(); stats = {"tasks": 0, "skipped": 0, "sub": 0, "done": 0, "err": 0} out = open(dst, "w") def work(tc): try: r = annotate(tc) except Exception: with lock: stats["err"] += 1 return if r is None: with lock: stats["skipped"] += 1 return plan, a = r rows = [] for i, c in enumerate(tc): rows.append({**c, "goal": plan[a[i]], "skill": "m2w@sub", "plan": plan, "sub_index": a[i]}) if EMIT_DONE and i > 0 and a[i] > a[i - 1]: rows.append({**c, "goal": plan[a[i - 1]], "gold_op": "DONE", "gold_id": "DONE", "kind": "done", "label": "", "skill": "subgoal_done", "plan": plan, "sub_index": a[i - 1]}) with lock: for r_ in rows: out.write(json.dumps(r_, ensure_ascii=False) + "\n") stats["tasks"] += 1; stats["sub"] += sum(r_["skill"] == "m2w@sub" for r_ in rows); stats["done"] += sum(r_["skill"] == "subgoal_done" for r_ in rows) if stats["tasks"] % 50 == 0: print(" ", stats, flush=True) with ThreadPoolExecutor(workers) as ex: list(ex.map(work, tasks)) out.close() print("done", stats, "->", dst) if __name__ == "__main__": main()