Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
(`<queue>__<n>.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).
Expand Down
11 changes: 10 additions & 1 deletion src/wfpy/_run_artifacts.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
2 changes: 2 additions & 0 deletions src/wfpy/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
53 changes: 53 additions & 0 deletions tests/test_queue_trace.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]]
Loading