Skip to content

Commit d73d54d

Browse files
authored
Artifact asset transfers (#2)
* Make CLI execution result output opt-in * Add artifact downloads and asset uploads
1 parent 239b097 commit d73d54d

10 files changed

Lines changed: 426 additions & 39 deletions

File tree

‎README.md‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -47,8 +47,8 @@ subfork execute GRAPH_ID --version v1
4747
subfork execute GRAPH_ID -o results.json
4848
```
4949

50-
`execute` waits for completion and returns JSON. Use `-o` to save results to a
51-
file; progress stays on stderr. Add `-f` / `--force` to overwrite existing output
50+
`execute` waits for completion and shows progress on stderr. Results are opt-in:
51+
use `-o results.json` to save JSON or `-o -` to print it to stdout. Add `-f` / `--force` to overwrite existing output
5252
files with `execute` or `export`. Run `subfork --help` or `subfork execute --help`
5353
for more options.
5454

@@ -120,7 +120,8 @@ remote state before repeating a create, publish, or execute request.
120120

121121
## Documentation
122122

123-
See the [documentation](docs/index.md) for installation, Python usage, CLI options,
123+
See the [API and file workflow guide](docs/api.md) for asset uploads and artifact
124+
downloads, and the [documentation](docs/index.md) for installation, Python usage, CLI options,
124125
and examples. Browse [Subfork Examples](https://examples.subfork.com) and the
125126
[subfork-examples repository](https://github.com/subforkdev/subfork-examples) for
126127
graphs to learn from and reuse.

‎docs/api.md‎

Lines changed: 119 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,119 @@
1+
# API reference and file workflows
2+
3+
The synchronous Python client returns dictionaries from the Subfork API. It does
4+
not print results; the command-line tool's output options do not affect Python.
5+
Set `SUBFORK_API_KEY`, then use a context manager to close connections reliably.
6+
7+
## Execute and inspect results
8+
9+
```python
10+
from typing import Any, Dict
11+
from subfork import Subfork
12+
13+
14+
def run_graph(client: Subfork, graph_id: str) -> Dict[str, Any]:
15+
"""Run a graph and return its completed outputs."""
16+
submitted = client.graphs.execute(graph_id)
17+
result = client.executions.wait(submitted["id"], timeout=300)
18+
if result["status"] != "completed":
19+
raise RuntimeError(f"Execution {submitted['id']}: {result['status']}")
20+
return result["outputs"]
21+
22+
23+
with Subfork() as client:
24+
outputs = run_graph(client, "g_YOUR_GRAPH")
25+
```
26+
27+
Output names and shapes are defined by each graph. Text and JSON can be used
28+
immediately. A media output may contain an `artifact` object, a remote URL, or
29+
inline data. Artifact download helpers accept an artifact ID, not arbitrary URLs.
30+
Waiting requires `graphs:read`; starting a run requires `graphs:run`. A waiting
31+
timeout does not cancel the remote execution.
32+
33+
## Download artifacts
34+
35+
For a graph exposing an `audio` media output with an artifact reference:
36+
37+
```python
38+
with Subfork() as client:
39+
outputs = run_graph(client, "g_YOUR_GRAPH")
40+
artifact_id = outputs["audio"]["artifact"]["artifact_id"]
41+
path = client.artifacts.download(artifact_id, "narration.mp3")
42+
```
43+
44+
A direct artifact output instead exposes `outputs["artifact"]["artifact_id"]`.
45+
The helper works for PDFs, images, MP3, MP4 and other stored file types. It returns
46+
`pathlib.Path`; it does not decode media or open a viewer. Use your preferred media
47+
library or player after downloading.
48+
49+
`download(artifact_id, destination, *, overwrite=False, max_bytes=250_000_000)`
50+
streams to a temporary file, then publishes the complete file at the requested
51+
path. The parent directory must exist. Existing files are protected unless you
52+
pass `overwrite=True`; failed downloads preserve them. The default decoded-byte
53+
limit is 250 MB and can be increased explicitly.
54+
55+
Downloads require `graphs:read` and access to the artifact. Expired or deleted
56+
artifacts cannot be recovered by the client. One HTTPS storage redirect is
57+
supported; the Subfork authorization header and cookies are not forwarded.
58+
59+
## Upload graph assets
60+
61+
**Server requirement:** API-key asset upload/list support must be deployed.
62+
Older servers return HTTP 403 even with an otherwise valid key.
63+
64+
```python
65+
with Subfork() as client:
66+
uploaded = client.assets.upload(
67+
"g_YOUR_GRAPH", "document.pdf", content_type="application/pdf"
68+
)
69+
artifact_id = uploaded["artifact_id"]
70+
print(artifact_id)
71+
```
72+
73+
`upload(graph_id, source, *, content_type=None)` streams a local file as multipart
74+
form data. The MIME type is inferred from its filename unless supplied; the
75+
server validates the file and controls persisted metadata. The returned JSON
76+
includes `artifact_id`, `media_type`, `size_bytes` and other artifact metadata.
77+
78+
Uploads require `graphs:write`, an owned graph, enabled uploads and available
79+
storage quota. The asset starts private. Uploading does **not** automatically
80+
bind an Asset node or execute the graph. To select it on an existing Asset node:
81+
82+
```python
83+
with Subfork() as client:
84+
graph_id = "g_YOUR_GRAPH"
85+
graph = client.graphs.get(graph_id)
86+
definition = graph["definition"]
87+
source_node = next(
88+
node for node in definition["nodes"]
89+
if node["node_instance_id"] == "pdf" and node["node_id"] == "n_asset"
90+
)
91+
uploaded = client.assets.upload(graph_id, "document.pdf")
92+
source_node["params"]["artifact_id"] = uploaded["artifact_id"]
93+
client.graphs.update(
94+
graph_id, name=graph["name"], definition=definition,
95+
description=graph.get("description") or "",
96+
)
97+
```
98+
99+
The example changes the draft; existing published versions remain immutable.
100+
Fetch the latest draft and avoid concurrent edits when updating its definition.
101+
For a graph exposing an artifact input, you can instead supply the reference in
102+
`graphs.execute(..., inputs={"pdf": {"artifact_id": artifact_id}})`. The input
103+
name must match that graph's published or draft interface.
104+
105+
List assets with `client.assets.list(graph_id)`. It returns an `assets` array in a
106+
JSON object. Use `include_generated=True` to include unexpired execution outputs.
107+
Listing requires `graphs:read` and ownership of the graph.
108+
109+
The client does not expose secret management, asset deletion, visibility changes
110+
or retention changes. Configure provider secrets through the Subfork UI.
111+
112+
## Errors and retries
113+
114+
HTTP failures raise `APIError` subclasses, including permission, validation and
115+
rate-limit errors. Downloads may also raise `FileExistsError`, filesystem errors,
116+
or `ValueError` for an invalid limit or oversized response. Network failures raise
117+
`TransportError`. File transfers and graph executions are never automatically
118+
retried: an upload or run may have succeeded before the connection failed. Inspect
119+
remote state before repeating a write.

‎docs/cli.md‎

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -38,21 +38,21 @@ subfork execute GRAPH_ID -o results.json
3838
subfork execute GRAPH_ID -o results.json --force
3939
```
4040

41-
Execution waits for completion and prints JSON outputs. `-o` / `--out` writes
42-
results to a file instead of stdout. Use `-f` / `--force` to overwrite an existing
41+
Execution waits for completion without printing results. `-o` / `--out` writes
42+
result JSON to a file; use `-o -` to print it to stdout. Use `-f` / `--force` to overwrite an existing
4343
file; this flag also works with `export`.
4444

4545
A yellow spinner shows active nodes. Finished nodes remain on stderr with their
46-
final status. JSON results stay on stdout, so piping works normally:
46+
final status. Opt in to JSON output for piping:
4747

4848
```bash
49-
subfork execute GRAPH_ID > results.json
49+
subfork execute GRAPH_ID -o - > results.json
5050
```
5151

5252
| Option | Behavior |
5353
| --- | --- |
54-
| `--no-wait` | Return the submission ID and status immediately |
55-
| `--raw` | Return the full execution snapshot |
54+
| `--no-wait` | Report submission ID and status on stderr immediately; use `-o` for JSON |
55+
| `--raw` | Print the full execution snapshot, or write it to the `-o` destination |
5656
| `--wait-timeout 300` | Wait up to 300 seconds; the default is 120 |
5757
| `--poll-interval 1` | Poll every second; the default is 2 |
5858

‎docs/mkpages.yml‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@ navigation:
88
href: /installation/
99
- label: Python
1010
href: /usage/
11+
- label: API
12+
href: /api/
1113
- label: CLI
1214
href: /cli/
1315
- label: Examples

‎docs/usage.md‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -87,3 +87,6 @@ requests. After a `TransportError`, check remote state before submitting another
8787
create, publish, or execute request.
8888

8989
See [troubleshooting](troubleshooting.md) for common errors.
90+
91+
See the [API and file workflow guide](api.md) for asset uploads, artifact downloads,
92+
and using execution results in your own functions.

‎src/subfork/cli.py‎

Lines changed: 30 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -183,7 +183,11 @@ def build_parser() -> argparse.ArgumentParser:
183183
child.add_argument("--comment", default="")
184184
elif command == "execute":
185185
child.add_argument("--version", default="draft")
186-
child.add_argument("-o", "--out", help="Write result JSON to a file instead of stdout")
186+
child.add_argument(
187+
"-o",
188+
"--out",
189+
help="Write result JSON to a file, or - for stdout (default: no results)",
190+
)
187191
child.add_argument(
188192
"--no-wait", action="store_true", help="Return submission status immediately"
189193
)
@@ -314,7 +318,12 @@ def main(argv: Optional[Sequence[str]] = None) -> int:
314318
if result_file is not None
315319
else (export_file if export_file != "-" else None)
316320
)
317-
if destination is not None and Path(destination).exists() and not force:
321+
if (
322+
destination is not None
323+
and destination != "-"
324+
and Path(destination).exists()
325+
and not force
326+
):
318327
raise ValueError("Output file already exists; use --force to overwrite.")
319328
with Subfork(base_url=args.base_url, timeout=args.timeout) as client:
320329
result = graph_command(client, args)
@@ -333,16 +342,25 @@ def main(argv: Optional[Sequence[str]] = None) -> int:
333342
}
334343
else:
335344
display = result.get("outputs", {})
336-
rendered = json.dumps(display, indent=2, ensure_ascii=False) + "\n"
337-
output = getattr(args, "output", "-")
338-
if result_file is None and output == "-":
339-
sys.stdout.write(rendered)
340-
else:
341-
# Open only after the request and serialization succeed.
342-
with Path(result_file if result_file is not None else output).open(
343-
"w" if force else "x", encoding="utf-8", newline="\n"
344-
) as stream:
345-
stream.write(rendered)
345+
emit_results = args.command != "execute" or result_file is not None or args.raw
346+
if emit_results:
347+
rendered = json.dumps(display, indent=2, ensure_ascii=False) + "\n"
348+
output = result_file if result_file is not None else getattr(args, "output", "-")
349+
if output == "-":
350+
sys.stdout.write(rendered)
351+
else:
352+
# Open only after the request and serialization succeed.
353+
with Path(output).open(
354+
"w" if force else "x", encoding="utf-8", newline="\n"
355+
) as stream:
356+
stream.write(rendered)
357+
elif args.command == "execute" and args.no_wait:
358+
print(
359+
"Execution {}: {}".format(
360+
result.get("id") or result.get("execution_id"), result.get("status")
361+
),
362+
file=sys.stderr,
363+
)
346364
if failed:
347365
print(
348366
"subfork: execution did not complete successfully; use --raw for details",

‎src/subfork/client.py‎

Lines changed: 71 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,8 @@
55
import math
66
import os
77
import re
8+
import tempfile
9+
from pathlib import Path
810
from types import TracebackType
911
from typing import Any, Optional, Type
1012

@@ -20,7 +22,7 @@
2022
TransportError,
2123
ValidationError,
2224
)
23-
from .resources import Executions, Graphs, Nodes
25+
from .resources import Artifacts, Assets, Executions, Graphs, Nodes
2426

2527

2628
def _error_message(response: httpx.Response) -> str:
@@ -30,7 +32,7 @@ def _error_message(response: httpx.Response) -> str:
3032
return message
3133
try:
3234
payload = response.json()
33-
except ValueError:
35+
except (ValueError, httpx.ResponseNotRead):
3436
return message
3537
detail = payload.get("detail") if isinstance(payload, dict) else None
3638
if isinstance(detail, str) and re.fullmatch(
@@ -107,6 +109,8 @@ def __init__(
107109
follow_redirects=False,
108110
transport=transport,
109111
)
112+
self.artifacts = Artifacts(self)
113+
self.assets = Assets(self)
110114
self.nodes = Nodes(self)
111115
self.graphs = Graphs(self)
112116
self.executions = Executions(self)
@@ -124,6 +128,17 @@ def _request(self, method: str, path: str, **kwargs: Any) -> Any:
124128
raise TransportError(
125129
"API request could not be completed; inspect remote state before retrying a write."
126130
) from None
131+
self._check_response(response)
132+
if response.status_code == 204:
133+
return None
134+
try:
135+
return response.json()
136+
except ValueError:
137+
raise InvalidResponseError("API returned invalid JSON.") from None
138+
139+
@staticmethod
140+
def _check_response(response: httpx.Response) -> None:
141+
"""Raise a credential-safe API exception for an unsuccessful response."""
127142
if not response.is_success:
128143
errors = {
129144
401: AuthenticationError,
@@ -138,12 +153,61 @@ def _request(self, method: str, path: str, **kwargs: Any) -> Any:
138153
status_code=response.status_code,
139154
retry_after=response.headers.get("retry-after"),
140155
)
141-
if response.status_code == 204:
142-
return None
156+
157+
def _download(self, path: str, destination: Path, *, overwrite: bool, max_bytes: int) -> Path:
158+
"""Stream an artifact, stripping credentials on a storage redirect."""
159+
if max_bytes <= 0:
160+
raise ValueError("max_bytes must be positive.")
161+
if destination.exists() and not overwrite:
162+
raise FileExistsError("Output file already exists; set overwrite=True.")
163+
temporary = None
164+
response = None
143165
try:
144-
return response.json()
145-
except ValueError:
146-
raise InvalidResponseError("API returned invalid JSON.") from None
166+
request = self._client.build_request("GET", path.lstrip("/"))
167+
response = self._client.send(request, stream=True)
168+
if response.status_code in {301, 302, 303, 307, 308}:
169+
location = response.headers.get("location")
170+
if not location:
171+
raise InvalidResponseError("Artifact redirect has no location.")
172+
try:
173+
target = response.url.join(location)
174+
except httpx.InvalidURL:
175+
raise InvalidResponseError(
176+
"Artifact redirect contains an invalid URL."
177+
) from None
178+
if target.scheme != "https" or target.userinfo or target.fragment:
179+
raise InvalidResponseError(
180+
"Artifact redirect must use HTTPS without credentials."
181+
)
182+
response.close()
183+
# A fresh Request does not inherit the API client's headers or cookies.
184+
response = self._client.send(
185+
httpx.Request("GET", target), stream=True, auth=None, follow_redirects=False
186+
)
187+
self._check_response(response)
188+
with tempfile.NamedTemporaryFile(
189+
dir=destination.parent, prefix=".subfork-", delete=False
190+
) as stream:
191+
temporary = Path(stream.name)
192+
size = 0
193+
for chunk in response.iter_bytes(chunk_size=65536):
194+
size += len(chunk)
195+
if size > max_bytes:
196+
raise ValueError("Artifact exceeds max_bytes; partial download discarded.")
197+
stream.write(chunk)
198+
if overwrite:
199+
os.replace(temporary, destination)
200+
else:
201+
# Exclusive publication also protects against another writer racing us.
202+
os.link(temporary, destination)
203+
return destination
204+
except httpx.TransportError:
205+
raise TransportError("Artifact download could not be completed.") from None
206+
finally:
207+
if response is not None:
208+
response.close()
209+
if temporary is not None:
210+
temporary.unlink(missing_ok=True)
147211

148212
def close(self) -> None:
149213
"""Release the HTTP connection pool."""

0 commit comments

Comments
 (0)