From 6ce033761d090d0ad4edf75081af6d3fdfa8bfd2 Mon Sep 17 00:00:00 2001 From: Endri Bezati Date: Fri, 2 Oct 2026 09:46:33 +0200 Subject: [PATCH] trace: each edge token gets its own file, so an earlier step opens the token then A token that is not a path is written to `work/edge-tokens/` for the overlay and the queue trace to point at. There was one file per queue, rewritten by every token, so every step's last token on an edge pointed at the same file -- and opening it at an earlier step in the IDE's stepper showed the run's final token. Each token now gets its own file, numbered per queue (`__.json`). Claude-Session: https://claude.ai/code/session_015VK7fH1c4aKbexnq2QcuKU --- CHANGELOG.md | 4 +++ src/wfpy/_run_artifacts.py | 11 +++++++- src/wfpy/runner.py | 2 ++ tests/test_queue_trace.py | 53 ++++++++++++++++++++++++++++++++++++++ 4 files changed, 69 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0e342f9..c95679f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -57,6 +57,10 @@ All notable changes to wfpy are documented here. The format follows firing. ### Fixed +- Each token that goes to `work/edge-tokens/` gets its own file + (`__.json`), instead of one file per queue rewritten by every + token. The queue trace's last token at an earlier step, which the IDE's + stepper opens, read as the run's final token. - A run whose tokens are arbitrary objects no longer fails writing its overlay or run record: an object those files cannot otherwise show is written as its `repr` (a dataclass as its fields). diff --git a/src/wfpy/_run_artifacts.py b/src/wfpy/_run_artifacts.py index 6782527..dbee9ae 100644 --- a/src/wfpy/_run_artifacts.py +++ b/src/wfpy/_run_artifacts.py @@ -67,7 +67,16 @@ def _materialize_edge_token(plan: Any, q: Any, value: Any) -> Any: token_dir = Path(work_dir) / "edge-tokens" token_dir.mkdir(parents=True, exist_ok=True) safe_queue_id = re.sub(r"[^A-Za-z0-9._-]+", "_", q.id).strip("_") or "queue" - token_path = token_dir / f"{safe_queue_id}.json" + # One file per token, numbered per queue. One file per queue was rewritten + # by every token, so the queue trace's last token at an early step -- what + # the IDE's stepper opens -- read as the run's final one. + counts = getattr(plan, "edge_token_counts", None) + if counts is None: + counts = {} + plan.edge_token_counts = counts + index = counts.get(q.id, 0) + counts[q.id] = index + 1 + token_path = token_dir / f"{safe_queue_id}__{index}.json" token_path.write_text(_runtime_json_dumps(value)) return str(token_path) diff --git a/src/wfpy/runner.py b/src/wfpy/runner.py index 4814886..fb0301d 100644 --- a/src/wfpy/runner.py +++ b/src/wfpy/runner.py @@ -513,6 +513,8 @@ def __init__(self, name: str = "") -> None: self.edge_info_by_queue_id: dict[str, dict[str, str]] = {} # queue_id → last value that traversed the edge self.edge_last_token_by_queue_id: dict[str, Any] = {} + # How many tokens each queue has written to edge-tokens/, numbering the files. + self.edge_token_counts: dict[str, int] = {} # Viewer overlay writer (set by run(), None in standalone execute_plan) self.overlay_writer: _ViewerOverlayWriter | None = None diff --git a/tests/test_queue_trace.py b/tests/test_queue_trace.py index 3a9f017..ed892a7 100644 --- a/tests/test_queue_trace.py +++ b/tests/test_queue_trace.py @@ -77,3 +77,56 @@ def test_a_queue_no_token_has_reached_carries_no_last_token(tmp_path: Path) -> N assert steps[0]["actorInstanceName"] == "c" assert "lastToken" not in first["d.Out-->WF.Out"] assert "lastToken" in first["c.Out-->d.In"] + + +@task +class Pairs: + """Emits a different list three times: tokens that go to edge-tokens/ files. + + Not a dict: a dict returned for one output is read as ``{port: value}``. + """ + + _n: int = 0 + + class Ports: + Out = Port[list](direction="out") + + @action(consumes={}, produces={"Out": 1}) + @guard(lambda self: self._n < 3) + def emit(self) -> list[int]: + self._n += 1 + return [self._n] + + +@task +class Keep: + class Ports: + In = Port[list](direction="in") + Out = Port[list](direction="out") + + @action(consumes={"In": 1}, produces={"Out": 1}) + def go(self, x: list[int]) -> list[int]: + return x + + +@workflow(outputs={"Out": list}) +def pairs(): + p = Pairs() + k = Keep() + connect(p.Out, k.In) + connect(k.Out, "Out") + + +def test_the_last_token_at_each_step_is_the_token_then(tmp_path: Path) -> None: + run(pairs, out_dir=str(tmp_path), run_id="r", max_workers=1) + steps = json.loads((tmp_path / "r" / "run.wf-queues.json").read_text())["steps"] + + # What the stepper opens for the edge into Keep, step by step: each token + # as it was when it was the last one, not the run's final token. + seen = [] + for step in steps: + queue = next(q for q in step["queueSizes"] if q["queueId"] == "p.Out-->k.In") + content = json.loads(Path(queue["lastToken"]).read_text()) + if not seen or seen[-1] != content: + seen.append(content) + assert seen == [[1], [2], [3]]