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]]