Skip to content

Commit afd5a60

Browse files
committed
Retain node progress and support execution result files
1 parent d9e93ca commit afd5a60

3 files changed

Lines changed: 181 additions & 29 deletions

File tree

‎README.md‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@ subfork create graph.json --name "My new graph"
4646
subfork interface GRAPH_ID
4747
subfork publish GRAPH_ID --version v1
4848
subfork execute GRAPH_ID --version v1 --inputs inputs.json
49+
subfork execute GRAPH_ID -o results.json
4950
```
5051

5152
Export writes a graph's draft definition, which can be passed directly to create.
@@ -54,6 +55,8 @@ Descriptions, tags, and published versions are not included in the export.
5455
JSON input files must contain an object; use `-` to read from stdin.
5556
Create makes a public draft, publish creates an immutable version, and execute
5657
waits for completion and prints the graph outputs. Runs consume account quota.
58+
Use `-o results.json` or `--out results.json` to write result JSON to a new file
59+
instead of stdout. Progress still goes to stderr. Existing files are not overwritten.
5760
Use `execute --no-wait` for a submission summary or `execute --raw` for the full
5861
execution snapshot. Waiting requires read and run scopes; submission alone requires
5962
run scope. `--wait-timeout` (default 120) and `--poll-interval` (default 2) are in
@@ -62,7 +65,8 @@ A yellow spinner precedes `Graph <name> .......... Running`, with `Running` in
6265
green. As status snapshots arrive, the line shows the currently running node titles
6366
(or multiple titles for parallel nodes). Short-lived nodes may finish between polls.
6467
Set `NO_COLOR` to disable colors. Progress goes to stderr; stdout remains JSON.
65-
Redirected progress uses a single plain-text line.
68+
Finished nodes remain on separate stderr lines with their final status before
69+
the JSON results appear. Redirected progress uses plain text without animation.
6670

6771
Other commands include `get`, `published`, and `versions`. Use `--help` on any
6872
command. Node discovery and execution monitoring are available through the Python API.

‎src/subfork/cli.py‎

Lines changed: 87 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,14 @@ def execution_status(graph_name: str) -> Iterator[Callable[[Dict[str, Any]], Non
2323
label = "Graph " + graph_name
2424
status = "Running"
2525
titles: Dict[str, str] = {}
26+
reported: Dict[str, str] = {}
27+
name_width = len(label)
2628
stopped = threading.Event()
29+
completed_marker = "✓"
30+
try:
31+
completed_marker.encode(stream.encoding or "ascii")
32+
except UnicodeError:
33+
completed_marker = "+"
2734
frames = "⠋⠙⠹⠸⠼⠴⠦⠧⠇⠏"
2835
try:
2936
frames.encode(stream.encoding or "ascii")
@@ -32,7 +39,7 @@ def execution_status(graph_name: str) -> Iterator[Callable[[Dict[str, Any]], Non
3239

3340
def update(snapshot: Dict[str, Any]) -> None:
3441
"""Replace the display state with active node titles from a snapshot."""
35-
nonlocal label, status
42+
nonlocal label, status, name_width
3643
definition = snapshot.get("definition") or {}
3744
for node in definition.get("nodes", []):
3845
titles[node["node_instance_id"]] = node.get("title") or node["node_instance_id"]
@@ -41,6 +48,43 @@ def update(snapshot: Dict[str, Any]) -> None:
4148
titles.get(key, key) for key, node in nodes.items() if node.get("status") == "running"
4249
]
4350
with lock:
51+
name_width = max(
52+
[name_width]
53+
+ [len("Node " + title) for title in titles.values()]
54+
+ [len("Node " + key) for key in nodes if key not in titles]
55+
)
56+
for key, node in nodes.items():
57+
outcome = node.get("status")
58+
if outcome not in {
59+
"completed",
60+
"failed",
61+
"canceled",
62+
"cancelled",
63+
"skipped",
64+
"outcome_unknown",
65+
}:
66+
continue
67+
if reported.get(key) == outcome:
68+
continue
69+
reported[key] = outcome
70+
state = "Canceled" if outcome == "cancelled" else outcome.replace("_", " ").title()
71+
name = "Node " + titles.get(key, key)
72+
if terminal:
73+
stream.write(
74+
"\r\033[2K"
75+
+ format_line(
76+
completed_marker if outcome == "completed" else "-", name, state
77+
)
78+
+ "\n"
79+
)
80+
else:
81+
stream.write(
82+
format_line(
83+
completed_marker if outcome == "completed" else "-", name, state
84+
)
85+
+ "\n"
86+
)
87+
stream.flush()
4488
label = (
4589
("Node " if len(active) == 1 else "Nodes ") + ", ".join(active)
4690
if active
@@ -52,25 +96,30 @@ def update(snapshot: Dict[str, Any]) -> None:
5296
else str(snapshot.get("status", "running")).replace("_", " ").title()
5397
)
5498

55-
def render(index: int) -> None:
56-
"""Draw one width-limited line with a yellow spinner and green run status."""
57-
with lock:
58-
name = "".join(char if char.isprintable() else " " for char in label)
59-
state = "".join(char if char.isprintable() else " " for char in status)
99+
def format_line(frame: str, name: str, state: str) -> str:
100+
"""Format a width-limited progress row without splitting color escapes."""
101+
name = "".join(char if char.isprintable() else " " for char in name)
102+
state = "".join(char if char.isprintable() else " " for char in state)
60103
width = max(1, shutil.get_terminal_size().columns - 1)
61-
name = name[: max(0, width - len(state) - 7)]
62-
dots = "." * max(3, min(10, width - len(name) - len(state) - 4))
63-
frame = frames[index % len(frames)]
104+
# Reserve space for the longest status so different outcomes align too.
105+
status_column = min(name_width + 14, max(7, width - len("Outcome Unknown")))
106+
name = name[: max(0, status_column - 7)]
107+
dots = "." * max(3, status_column - len(name) - 4)
64108
plain = "{} {} {} {}".format(frame, name, dots, state)
65109
if len(plain) > width:
66-
line = plain[:width]
67-
elif os.environ.get("NO_COLOR"):
68-
line = plain
69-
else:
70-
colored_state = "\033[32m" + state + "\033[0m" if state == "Running" else state
71-
line = "\033[33m{}\033[0m {} {} {}".format(frame, name, dots, colored_state)
72-
stream.write("\r\033[2K" + line)
73-
stream.flush()
110+
return plain[:width]
111+
if not terminal or os.environ.get("NO_COLOR"):
112+
return plain
113+
color = "32" if state in {"Running", "Completed"} else "31" if state == "Failed" else "33"
114+
colored_state = "\033[" + color + "m" + state + "\033[0m"
115+
marker_color = "32" if state == "Completed" else "33"
116+
return "\033[{}m{}\033[0m {} {} {}".format(marker_color, frame, name, dots, colored_state)
117+
118+
def render(index: int) -> None:
119+
"""Draw the active row without interleaving retained completion lines."""
120+
with lock:
121+
stream.write("\r\033[2K" + format_line(frames[index % len(frames)], label, status))
122+
stream.flush()
74123

75124
def animate() -> None:
76125
"""Refresh the spinner independently of network requests and polling."""
@@ -122,6 +171,9 @@ def build_parser() -> argparse.ArgumentParser:
122171
child.add_argument("--comment", default="")
123172
elif command == "execute":
124173
child.add_argument("--version", default="draft")
174+
child.add_argument(
175+
"-o", "--out", help="Write result JSON to a new file instead of stdout"
176+
)
125177
child.add_argument(
126178
"--no-wait", action="store_true", help="Return submission status immediately"
127179
)
@@ -205,13 +257,14 @@ def graph_command(client: Subfork, args: argparse.Namespace) -> Any:
205257
):
206258
raise ValueError("Wait timeout and poll interval must be positive finite numbers.")
207259
result = graphs.execute(args.graph_id, version=args.version, inputs=inputs)
208-
if not args.no_wait and result.get("status") not in {
260+
terminal_statuses = {
209261
"completed",
210262
"failed",
211263
"canceled",
212264
"cancelled",
213265
"outcome_unknown",
214-
}:
266+
}
267+
if not args.no_wait:
215268
execution_id = result.get("id") or result.get("execution_id")
216269
if not isinstance(execution_id, str) or not execution_id:
217270
raise ValueError(
@@ -225,12 +278,13 @@ def graph_command(client: Subfork, args: argparse.Namespace) -> Any:
225278
graph_name if isinstance(graph_name, str) else args.graph_id
226279
) as update:
227280
update(result)
228-
result = client.executions.wait(
229-
execution_id,
230-
timeout=args.wait_timeout,
231-
poll_interval=args.poll_interval,
232-
on_update=update,
233-
)
281+
if result.get("status") not in terminal_statuses:
282+
result = client.executions.wait(
283+
execution_id,
284+
timeout=args.wait_timeout,
285+
poll_interval=args.poll_interval,
286+
on_update=update,
287+
)
234288
return result
235289

236290

@@ -242,6 +296,9 @@ def main(argv: Optional[Sequence[str]] = None) -> int:
242296
"""
243297
args = build_parser().parse_args(argv)
244298
try:
299+
result_file = getattr(args, "out", None)
300+
if result_file is not None and Path(result_file).exists():
301+
raise ValueError("Output file already exists; choose a new path.")
245302
with Subfork(base_url=args.base_url, timeout=args.timeout) as client:
246303
result = graph_command(client, args)
247304
failed = args.command == "execute" and result.get("status") in {
@@ -261,11 +318,13 @@ def main(argv: Optional[Sequence[str]] = None) -> int:
261318
display = result.get("outputs", {})
262319
rendered = json.dumps(display, indent=2, ensure_ascii=False) + "\n"
263320
output = getattr(args, "output", "-")
264-
if output == "-":
321+
if result_file is None and output == "-":
265322
sys.stdout.write(rendered)
266323
else:
267-
# Exclusive creation avoids silently overwriting a local definition.
268-
with Path(output).open("x", encoding="utf-8", newline="\n") as stream:
324+
# Exclusive creation avoids silently overwriting an existing file.
325+
with Path(result_file if result_file is not None else output).open(
326+
"x", encoding="utf-8", newline="\n"
327+
) as stream:
269328
stream.write(rendered)
270329
if failed:
271330
print(

‎tests/test_cli.py‎

Lines changed: 89 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22

33
import io
44
import json
5+
import re
56
from pathlib import Path
67
from typing import Any
78

@@ -260,3 +261,91 @@ def write(self, value: str) -> int:
260261
assert "\033[33m" in terminal.getvalue()
261262
assert "\033[32mRunning\033[0m" in terminal.getvalue()
262263
assert not any(thread.name == "subfork-spinner" for thread in threading.enumerate())
264+
265+
266+
@pytest.mark.parametrize("terminal", [False, True])
267+
def test_retained_node_outcomes(
268+
monkeypatch: pytest.MonkeyPatch, capsys: Any, terminal: bool
269+
) -> None:
270+
"""Retain each terminal node once, including nodes that finish between polls."""
271+
272+
class ProgressStream(io.StringIO):
273+
"""Capture either terminal or redirected progress."""
274+
275+
def isatty(self) -> bool:
276+
"""Select the requested output mode."""
277+
return terminal
278+
279+
stream = ProgressStream()
280+
monkeypatch.setattr("sys.stderr", stream)
281+
monkeypatch.setenv("TERM", "xterm")
282+
monkeypatch.setenv("NO_COLOR", "1")
283+
snapshot = {
284+
"status": "failed",
285+
"definition": {
286+
"nodes": [
287+
{"node_instance_id": "a", "title": "Fetch"},
288+
{"node_instance_id": "b", "title": "Parse"},
289+
{"node_instance_id": "c", "title": "Save"},
290+
]
291+
},
292+
"node_executions": {
293+
"a": {"status": "completed"},
294+
"b": {"status": "failed"},
295+
"c": {"status": "canceled"},
296+
},
297+
}
298+
with cli.execution_status("Pipeline") as update:
299+
update(snapshot)
300+
update(snapshot)
301+
output = stream.getvalue()
302+
for name, outcome in (("Fetch", "Completed"), ("Parse", "Failed"), ("Save", "Canceled")):
303+
assert len(re.findall(r"Node " + name + r" \.+ " + outcome + r"\n", output)) == 1
304+
assert capsys.readouterr().out == ""
305+
306+
307+
def test_node_status_columns_align(monkeypatch: pytest.MonkeyPatch) -> None:
308+
"""Pad different node names so terminal outcomes start in the same column."""
309+
stream = io.StringIO()
310+
monkeypatch.setattr("sys.stderr", stream)
311+
with cli.execution_status("Demo") as update:
312+
update(
313+
{
314+
"node_executions": {
315+
"A": {"status": "completed"},
316+
"Longer node name": {"status": "failed"},
317+
"Mid": {"status": "canceled"},
318+
}
319+
}
320+
)
321+
lines = stream.getvalue().splitlines()[1:]
322+
columns = [line.index(state) for line, state in zip(lines, ["Completed", "Failed", "Canceled"])]
323+
assert len(columns) == 3
324+
assert len(set(columns)) == 1
325+
326+
327+
@pytest.mark.parametrize("flag", ["-o", "--out"])
328+
@pytest.mark.parametrize("mode", [None, "--raw", "--no-wait"])
329+
def test_execute_result_file(
330+
requests: list, tmp_path: Path, capsys: Any, flag: str, mode: Any
331+
) -> None:
332+
"""Route results exclusively to a file and reject overwrites before submission."""
333+
output = tmp_path / "result.json"
334+
arguments = ["execute", "g_test", flag, str(output)]
335+
if mode:
336+
arguments.append(mode)
337+
assert cli.main(arguments) == 0
338+
value = json.loads(output.read_text())
339+
assert value == (
340+
{"id": "test", "status": "completed"}
341+
if mode == "--raw"
342+
else {"execution_id": "test", "status": "completed"} if mode == "--no-wait" else {}
343+
)
344+
assert capsys.readouterr().out == ""
345+
count = len(requests)
346+
assert cli.main(arguments) == 1
347+
assert len(requests) == count
348+
assert json.loads(output.read_text()) == value
349+
captured = capsys.readouterr()
350+
assert captured.out == ""
351+
assert "already exists" in captured.err

0 commit comments

Comments
 (0)