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
7 changes: 4 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -47,8 +47,8 @@ subfork execute GRAPH_ID --version v1
subfork execute GRAPH_ID -o results.json
```

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

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

## Documentation

See the [documentation](docs/index.md) for installation, Python usage, CLI options,
See the [API and file workflow guide](docs/api.md) for asset uploads and artifact
downloads, and the [documentation](docs/index.md) for installation, Python usage, CLI options,
and examples. Browse [Subfork Examples](https://examples.subfork.com) and the
[subfork-examples repository](https://github.com/subforkdev/subfork-examples) for
graphs to learn from and reuse.
Expand Down
119 changes: 119 additions & 0 deletions docs/api.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
# API reference and file workflows

The synchronous Python client returns dictionaries from the Subfork API. It does
not print results; the command-line tool's output options do not affect Python.
Set `SUBFORK_API_KEY`, then use a context manager to close connections reliably.

## Execute and inspect results

```python
from typing import Any, Dict
from subfork import Subfork


def run_graph(client: Subfork, graph_id: str) -> Dict[str, Any]:
"""Run a graph and return its completed outputs."""
submitted = client.graphs.execute(graph_id)
result = client.executions.wait(submitted["id"], timeout=300)
if result["status"] != "completed":
raise RuntimeError(f"Execution {submitted['id']}: {result['status']}")
return result["outputs"]


with Subfork() as client:
outputs = run_graph(client, "g_YOUR_GRAPH")
```

Output names and shapes are defined by each graph. Text and JSON can be used
immediately. A media output may contain an `artifact` object, a remote URL, or
inline data. Artifact download helpers accept an artifact ID, not arbitrary URLs.
Waiting requires `graphs:read`; starting a run requires `graphs:run`. A waiting
timeout does not cancel the remote execution.

## Download artifacts

For a graph exposing an `audio` media output with an artifact reference:

```python
with Subfork() as client:
outputs = run_graph(client, "g_YOUR_GRAPH")
artifact_id = outputs["audio"]["artifact"]["artifact_id"]
path = client.artifacts.download(artifact_id, "narration.mp3")
```

A direct artifact output instead exposes `outputs["artifact"]["artifact_id"]`.
The helper works for PDFs, images, MP3, MP4 and other stored file types. It returns
`pathlib.Path`; it does not decode media or open a viewer. Use your preferred media
library or player after downloading.

`download(artifact_id, destination, *, overwrite=False, max_bytes=250_000_000)`
streams to a temporary file, then publishes the complete file at the requested
path. The parent directory must exist. Existing files are protected unless you
pass `overwrite=True`; failed downloads preserve them. The default decoded-byte
limit is 250 MB and can be increased explicitly.

Downloads require `graphs:read` and access to the artifact. Expired or deleted
artifacts cannot be recovered by the client. One HTTPS storage redirect is
supported; the Subfork authorization header and cookies are not forwarded.

## Upload graph assets

**Server requirement:** API-key asset upload/list support must be deployed.
Older servers return HTTP 403 even with an otherwise valid key.

```python
with Subfork() as client:
uploaded = client.assets.upload(
"g_YOUR_GRAPH", "document.pdf", content_type="application/pdf"
)
artifact_id = uploaded["artifact_id"]
print(artifact_id)
```

`upload(graph_id, source, *, content_type=None)` streams a local file as multipart
form data. The MIME type is inferred from its filename unless supplied; the
server validates the file and controls persisted metadata. The returned JSON
includes `artifact_id`, `media_type`, `size_bytes` and other artifact metadata.

Uploads require `graphs:write`, an owned graph, enabled uploads and available
storage quota. The asset starts private. Uploading does **not** automatically
bind an Asset node or execute the graph. To select it on an existing Asset node:

```python
with Subfork() as client:
graph_id = "g_YOUR_GRAPH"
graph = client.graphs.get(graph_id)
definition = graph["definition"]
source_node = next(
node for node in definition["nodes"]
if node["node_instance_id"] == "pdf" and node["node_id"] == "n_asset"
)
uploaded = client.assets.upload(graph_id, "document.pdf")
source_node["params"]["artifact_id"] = uploaded["artifact_id"]
client.graphs.update(
graph_id, name=graph["name"], definition=definition,
description=graph.get("description") or "",
)
```

The example changes the draft; existing published versions remain immutable.
Fetch the latest draft and avoid concurrent edits when updating its definition.
For a graph exposing an artifact input, you can instead supply the reference in
`graphs.execute(..., inputs={"pdf": {"artifact_id": artifact_id}})`. The input
name must match that graph's published or draft interface.

List assets with `client.assets.list(graph_id)`. It returns an `assets` array in a
JSON object. Use `include_generated=True` to include unexpired execution outputs.
Listing requires `graphs:read` and ownership of the graph.

The client does not expose secret management, asset deletion, visibility changes
or retention changes. Configure provider secrets through the Subfork UI.

## Errors and retries

HTTP failures raise `APIError` subclasses, including permission, validation and
rate-limit errors. Downloads may also raise `FileExistsError`, filesystem errors,
or `ValueError` for an invalid limit or oversized response. Network failures raise
`TransportError`. File transfers and graph executions are never automatically
retried: an upload or run may have succeeded before the connection failed. Inspect
remote state before repeating a write.
12 changes: 6 additions & 6 deletions docs/cli.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,21 +38,21 @@ subfork execute GRAPH_ID -o results.json
subfork execute GRAPH_ID -o results.json --force
```

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

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

```bash
subfork execute GRAPH_ID > results.json
subfork execute GRAPH_ID -o - > results.json
```

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

Expand Down
2 changes: 2 additions & 0 deletions docs/mkpages.yml
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ navigation:
href: /installation/
- label: Python
href: /usage/
- label: API
href: /api/
- label: CLI
href: /cli/
- label: Examples
Expand Down
3 changes: 3 additions & 0 deletions docs/usage.md
Original file line number Diff line number Diff line change
Expand Up @@ -87,3 +87,6 @@ requests. After a `TransportError`, check remote state before submitting another
create, publish, or execute request.

See [troubleshooting](troubleshooting.md) for common errors.

See the [API and file workflow guide](api.md) for asset uploads, artifact downloads,
and using execution results in your own functions.
42 changes: 30 additions & 12 deletions src/subfork/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -183,7 +183,11 @@ def build_parser() -> argparse.ArgumentParser:
child.add_argument("--comment", default="")
elif command == "execute":
child.add_argument("--version", default="draft")
child.add_argument("-o", "--out", help="Write result JSON to a file instead of stdout")
child.add_argument(
"-o",
"--out",
help="Write result JSON to a file, or - for stdout (default: no results)",
)
child.add_argument(
"--no-wait", action="store_true", help="Return submission status immediately"
)
Expand Down Expand Up @@ -314,7 +318,12 @@ def main(argv: Optional[Sequence[str]] = None) -> int:
if result_file is not None
else (export_file if export_file != "-" else None)
)
if destination is not None and Path(destination).exists() and not force:
if (
destination is not None
and destination != "-"
and Path(destination).exists()
and not force
):
raise ValueError("Output file already exists; use --force to overwrite.")
with Subfork(base_url=args.base_url, timeout=args.timeout) as client:
result = graph_command(client, args)
Expand All @@ -333,16 +342,25 @@ def main(argv: Optional[Sequence[str]] = None) -> int:
}
else:
display = result.get("outputs", {})
rendered = json.dumps(display, indent=2, ensure_ascii=False) + "\n"
output = getattr(args, "output", "-")
if result_file is None and output == "-":
sys.stdout.write(rendered)
else:
# Open only after the request and serialization succeed.
with Path(result_file if result_file is not None else output).open(
"w" if force else "x", encoding="utf-8", newline="\n"
) as stream:
stream.write(rendered)
emit_results = args.command != "execute" or result_file is not None or args.raw
if emit_results:
rendered = json.dumps(display, indent=2, ensure_ascii=False) + "\n"
output = result_file if result_file is not None else getattr(args, "output", "-")
if output == "-":
sys.stdout.write(rendered)
else:
# Open only after the request and serialization succeed.
with Path(output).open(
"w" if force else "x", encoding="utf-8", newline="\n"
) as stream:
stream.write(rendered)
elif args.command == "execute" and args.no_wait:
print(
"Execution {}: {}".format(
result.get("id") or result.get("execution_id"), result.get("status")
),
file=sys.stderr,
)
if failed:
print(
"subfork: execution did not complete successfully; use --raw for details",
Expand Down
78 changes: 71 additions & 7 deletions src/subfork/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@
import math
import os
import re
import tempfile
from pathlib import Path
from types import TracebackType
from typing import Any, Optional, Type

Expand All @@ -20,7 +22,7 @@
TransportError,
ValidationError,
)
from .resources import Executions, Graphs, Nodes
from .resources import Artifacts, Assets, Executions, Graphs, Nodes


def _error_message(response: httpx.Response) -> str:
Expand All @@ -30,7 +32,7 @@ def _error_message(response: httpx.Response) -> str:
return message
try:
payload = response.json()
except ValueError:
except (ValueError, httpx.ResponseNotRead):
return message
detail = payload.get("detail") if isinstance(payload, dict) else None
if isinstance(detail, str) and re.fullmatch(
Expand Down Expand Up @@ -107,6 +109,8 @@ def __init__(
follow_redirects=False,
transport=transport,
)
self.artifacts = Artifacts(self)
self.assets = Assets(self)
self.nodes = Nodes(self)
self.graphs = Graphs(self)
self.executions = Executions(self)
Expand All @@ -124,6 +128,17 @@ def _request(self, method: str, path: str, **kwargs: Any) -> Any:
raise TransportError(
"API request could not be completed; inspect remote state before retrying a write."
) from None
self._check_response(response)
if response.status_code == 204:
return None
try:
return response.json()
except ValueError:
raise InvalidResponseError("API returned invalid JSON.") from None

@staticmethod
def _check_response(response: httpx.Response) -> None:
"""Raise a credential-safe API exception for an unsuccessful response."""
if not response.is_success:
errors = {
401: AuthenticationError,
Expand All @@ -138,12 +153,61 @@ def _request(self, method: str, path: str, **kwargs: Any) -> Any:
status_code=response.status_code,
retry_after=response.headers.get("retry-after"),
)
if response.status_code == 204:
return None

def _download(self, path: str, destination: Path, *, overwrite: bool, max_bytes: int) -> Path:
"""Stream an artifact, stripping credentials on a storage redirect."""
if max_bytes <= 0:
raise ValueError("max_bytes must be positive.")
if destination.exists() and not overwrite:
raise FileExistsError("Output file already exists; set overwrite=True.")
temporary = None
response = None
try:
return response.json()
except ValueError:
raise InvalidResponseError("API returned invalid JSON.") from None
request = self._client.build_request("GET", path.lstrip("/"))
response = self._client.send(request, stream=True)
if response.status_code in {301, 302, 303, 307, 308}:
location = response.headers.get("location")
if not location:
raise InvalidResponseError("Artifact redirect has no location.")
try:
target = response.url.join(location)
except httpx.InvalidURL:
raise InvalidResponseError(
"Artifact redirect contains an invalid URL."
) from None
if target.scheme != "https" or target.userinfo or target.fragment:
raise InvalidResponseError(
"Artifact redirect must use HTTPS without credentials."
)
response.close()
# A fresh Request does not inherit the API client's headers or cookies.
response = self._client.send(
httpx.Request("GET", target), stream=True, auth=None, follow_redirects=False
)
self._check_response(response)
with tempfile.NamedTemporaryFile(
dir=destination.parent, prefix=".subfork-", delete=False
) as stream:
temporary = Path(stream.name)
size = 0
for chunk in response.iter_bytes(chunk_size=65536):
size += len(chunk)
if size > max_bytes:
raise ValueError("Artifact exceeds max_bytes; partial download discarded.")
stream.write(chunk)
if overwrite:
os.replace(temporary, destination)
else:
# Exclusive publication also protects against another writer racing us.
os.link(temporary, destination)
return destination
except httpx.TransportError:
raise TransportError("Artifact download could not be completed.") from None
finally:
if response is not None:
response.close()
if temporary is not None:
temporary.unlink(missing_ok=True)

def close(self) -> None:
"""Release the HTTP connection pool."""
Expand Down
Loading
Loading