diff --git a/.editorconfig b/.editorconfig
new file mode 100644
index 0000000..61312b9
--- /dev/null
+++ b/.editorconfig
@@ -0,0 +1,15 @@
+root = true
+
+[*]
+charset = utf-8
+end_of_line = lf
+insert_final_newline = true
+trim_trailing_whitespace = true
+indent_style = space
+indent_size = 4
+
+[*.{yml,yaml,toml}]
+indent_size = 2
+
+[Makefile]
+indent_style = tab
diff --git a/.flake8 b/.flake8
new file mode 100644
index 0000000..f210cfd
--- /dev/null
+++ b/.flake8
@@ -0,0 +1,7 @@
+[flake8]
+max-line-length = 100
+max-complexity = 35
+exclude = .git,__pycache__,.venv*,venv,build,dist,*.egg-info
+extend-ignore = E203,E501,W503
+# Enforce documentation for modules, classes, functions, methods, and constructors.
+docstring-convention = google
diff --git a/.gitattributes b/.gitattributes
new file mode 100644
index 0000000..e3e4309
--- /dev/null
+++ b/.gitattributes
@@ -0,0 +1,3 @@
+# Normalize text to LF in the repository and on checkout, on every platform.
+# Git detects binary files and leaves their contents unchanged.
+* text=auto eol=lf
diff --git a/.github/workflows/lint.yml b/.github/workflows/lint.yml
new file mode 100644
index 0000000..e9f0d3d
--- /dev/null
+++ b/.github/workflows/lint.yml
@@ -0,0 +1,18 @@
+name: Lint and type check
+
+on: [push, pull_request]
+
+permissions:
+ contents: read
+
+jobs:
+ lint:
+ runs-on: ubuntu-latest
+ steps:
+ - uses: actions/checkout@v4
+ - uses: actions/setup-python@v5
+ with:
+ python-version: '3.12'
+ - run: python -m pip install --upgrade pip
+ - run: python -m pip install -e '.[dev]'
+ - run: make lint typecheck
diff --git a/.github/workflows/pages.yml b/.github/workflows/pages.yml
new file mode 100644
index 0000000..7afb688
--- /dev/null
+++ b/.github/workflows/pages.yml
@@ -0,0 +1,59 @@
+name: Pages
+
+on:
+ push:
+ branches: [master]
+ pull_request:
+ workflow_dispatch:
+
+permissions:
+ contents: read
+
+concurrency:
+ group: pages-${{ github.ref }}
+ cancel-in-progress: true
+
+jobs:
+ build:
+ runs-on: ubuntu-latest
+ timeout-minutes: 10
+ steps:
+ - uses: actions/checkout@v4
+ - uses: actions/setup-python@v5
+ with:
+ python-version: "3.12"
+ cache: pip
+ cache-dependency-path: pyproject.toml
+ - name: Install MkPages
+ run: python -m pip install '.[docs]'
+ - name: Generate Jekyll source
+ run: >-
+ mkpages build docs --output .mkpages
+ --url "https://subforkdev.github.io"
+ --baseurl "/subfork-python"
+ - name: Build site
+ uses: actions/jekyll-build-pages@v1
+ with:
+ source: ./.mkpages
+ destination: ./_site
+ - name: Upload Pages artifact
+ uses: actions/upload-pages-artifact@v3
+ with:
+ path: ./_site
+
+ deploy:
+ if: github.event_name != 'pull_request' && github.ref == 'refs/heads/master'
+ needs: build
+ runs-on: ubuntu-latest
+ permissions:
+ pages: write
+ id-token: write
+ environment:
+ name: github-pages
+ url: ${{ steps.deployment.outputs.page_url }}
+ steps:
+ - name: Configure Pages
+ uses: actions/configure-pages@v5
+ - name: Deploy Pages
+ id: deployment
+ uses: actions/deploy-pages@v4
diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml
new file mode 100644
index 0000000..7b17a63
--- /dev/null
+++ b/.github/workflows/tests.yml
@@ -0,0 +1,20 @@
+name: Test and package
+on: [push, pull_request]
+permissions:
+ contents: read
+jobs:
+ test:
+ runs-on: ubuntu-latest
+ strategy:
+ matrix:
+ python: ['3.8', '3.10', '3.12', '3.13', '3.14']
+ steps:
+ - uses: actions/checkout@v4
+ - uses: actions/setup-python@v5
+ with:
+ python-version: ${{ matrix.python }}
+ - run: python -m pip install --upgrade pip
+ - run: python -m pip install '.[dev]'
+ - run: python -m pytest
+ - run: python -m build
+ - run: python -m twine check dist/*
diff --git a/.gitignore b/.gitignore
new file mode 100644
index 0000000..5a6ccce
--- /dev/null
+++ b/.gitignore
@@ -0,0 +1,56 @@
+# Python bytecode and compiled extensions
+__pycache__/
+*.py[cod]
+*$py.class
+*.so
+
+# Packaging and build output
+build/
+dist/
+.eggs/
+*.egg-info/
+*.egg
+MANIFEST
+
+# Virtual environments and local credentials
+.venv*/
+venv*/
+.env
+.env.*
+!.env.example
+!.env.*.example
+
+# Test, coverage, type-checker, and linter caches
+.pytest_cache/
+.mypy_cache/
+.ruff_cache/
+.tox/
+.nox/
+.coverage
+.coverage.*
+htmlcov/
+coverage.xml
+junit.xml
+.hypothesis/
+
+# Local scratch data and notebook checkpoints
+/tmp/
+/logs/
+.ipynb_checkpoints/
+
+# Editor, OS, and agent-local settings
+.vscode/
+.idea/
+*.swp
+*.swo
+*~
+.DS_Store
+Thumbs.db
+Desktop.ini
+.codex/
+.agents/
+
+# Generated documentation site
+.mkpages/
+_site/
+.jekyll-cache/
diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md
new file mode 100644
index 0000000..99fd037
--- /dev/null
+++ b/CONTRIBUTING.md
@@ -0,0 +1,67 @@
+# Contributing
+
+To contribute to this Python client, install its development dependencies using
+Python 3.8+ and a modern pip:
+
+```bash
+python -m pip install --upgrade pip
+python -m pip install -e '.[dev]'
+```
+
+## Code quality
+
+Install `.[dev]` to get the pinned development tools. Black 24.8.0, isort 5.13.2,
+and Flake8 7.1.1 use 100-column formatting, isort's Black profile, and a Python
+3.8 target.
+
+```bash
+make format # Sort imports and apply Black to src/ and tests/
+make lint # Check formatting, imports, Flake8 and Google-style docstrings
+make typecheck # Check annotated library code with mypy
+make test # Run the HTTP contract tests
+make check # Run lint, type checking and tests
+make build # Build distributions and check their metadata
+```
+
+Activate your environment first, or specify it explicitly, for example:
+`make check PYTHON=.venv/bin/python`. The CI lint job uses the same commands.
+EditorConfig defines whitespace and newline conventions for supporting editors.
+
+Add annotations and docstrings to new functions, methods and classes, including
+constructors and test helpers. Public docstrings should explain side effects,
+permissions, return values and failure behavior where useful. JSON dictionaries
+retain `Any` values because node parameters and server response fields are dynamic;
+this release does not pretend to provide complete generated response models.
+Mypy checks library signatures and bodies; Flake8 and formatting cover both the
+library and tests. Runtime tests continue to cover Python 3.8 in CI.
+
+## Versioning
+
+```bash
+make version # Show the current version
+make check-version # Verify both version fields agree
+make bump-patch # 2.0.0 -> 2.0.1
+make bump-minor # 2.0.0 -> 2.1.0
+make bump-major # 2.0.0 -> 3.0.0
+make bump-version VERSION=2.1.0 # Set an explicit stable version
+```
+
+Bumps update `pyproject.toml` and `subfork.__version__` together. They do not
+commit, tag, or publish. Review and commit the changes before releasing.
+
+## Documentation site
+
+Public documentation lives in `docs/`. Its `mkpages.yml` selects the dark theme.
+To generate the Jekyll source with Python 3.12:
+
+```bash
+python -m pip install -e '.[docs]'
+mkpages build docs --output .mkpages
+mkpages preview docs/
+```
+
+The Pages workflow builds pull requests and deploys `master`. In repository
+**Settings → Pages**, select **GitHub Actions** as the source. The initial site
+URL is `https://subforkdev.github.io/subfork-python/`; the workflow supplies its
+`--url` and `--baseurl` options. Update those if a custom domain is configured.
+Preview serves the documentation at `http://127.0.0.1:4000/`.
diff --git a/LICENSE b/LICENSE
new file mode 100644
index 0000000..5520dd3
--- /dev/null
+++ b/LICENSE
@@ -0,0 +1,29 @@
+BSD 3-Clause License
+
+Copyright (c) 2026, Ryan Galloway
+All rights reserved.
+
+Redistribution and use in source and binary forms, with or without
+modification, are permitted provided that the following conditions are met:
+
+1. Redistributions of source code must retain the above copyright notice, this
+ list of conditions and the following disclaimer.
+
+2. Redistributions in binary form must reproduce the above copyright notice,
+ this list of conditions and the following disclaimer in the documentation
+ and/or other materials provided with the distribution.
+
+3. Neither the name of the copyright holder nor the names of its
+ contributors may be used to endorse or promote products derived from
+ this software without specific prior written permission.
+
+THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS"
+AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE
+IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
+DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE
+FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL
+DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR
+SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER
+CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY,
+OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
+OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
diff --git a/Makefile b/Makefile
new file mode 100644
index 0000000..ef8ad4a
--- /dev/null
+++ b/Makefile
@@ -0,0 +1,38 @@
+PYTHON ?= python
+
+.PHONY: format lint typecheck test check build
+format:
+ $(PYTHON) -m isort src tests scripts
+ $(PYTHON) -m black src tests scripts
+
+lint:
+ $(PYTHON) -m black --check src tests scripts
+ $(PYTHON) -m isort --check-only src tests scripts
+ $(PYTHON) -m flake8 src tests scripts
+
+typecheck:
+ $(PYTHON) -m mypy
+
+test:
+ $(PYTHON) -m pytest
+
+check: lint typecheck test check-version
+
+build: check-version
+ $(PYTHON) -m build
+ $(PYTHON) -m twine check dist/*
+
+.PHONY: version check-version bump-patch bump-minor bump-major bump-version
+export VERSION
+
+version:
+ @$(PYTHON) scripts/bump-version.py --show
+
+check-version:
+ @$(PYTHON) scripts/bump-version.py --check
+
+bump-patch bump-minor bump-major:
+ @$(PYTHON) scripts/bump-version.py --bump $(@:bump-%=%)
+
+bump-version:
+ @$(PYTHON) scripts/bump-version.py --set
diff --git a/README.md b/README.md
index 3b18e51..79a0713 100644
--- a/README.md
+++ b/README.md
@@ -1 +1,134 @@
-hello world
+
+
+# Subfork Python
+
+The Python client for [Subfork](https://subfork.com). Discover nodes, build reusable
+graphs, publish versions, and run them from your scripts, applications, or agents.
+
+## Install
+
+Requires Python 3.8+.
+
+```bash
+pip install subfork
+```
+
+## Connect
+
+Create a key under **Account → API keys** in Subfork and set `SUBFORK_API_KEY` in
+your environment. Grant the permissions your application needs: read, write, run,
+and/or publish. Keep the key out of source files and graph definitions.
+
+```python
+from subfork import Subfork
+
+with Subfork() as client:
+ nodes = client.nodes.list()
+ graphs = client.graphs.list()
+```
+
+The client reads `SUBFORK_API_KEY` and connects to [subfork.com](https://subfork.com).
+
+## Command line
+
+The package includes a `subfork` CLI that uses your `SUBFORK_API_KEY`.
+
+```bash
+subfork list
+subfork export GRAPH_ID --output graph.json
+subfork validate graph.json
+subfork create graph.json --name "My new graph"
+subfork publish GRAPH_ID --version v1
+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
+files with `execute` or `export`. Run `subfork --help` or `subfork execute --help`
+for more options.
+
+## Create and run a graph
+
+This example creates a text-producing graph, publishes `v1`, and runs that version.
+It requires read, write, publish, and run permissions. Graphs are public, and runs
+use your account's execution quota.
+
+```python
+from subfork import Subfork
+
+name = "Python greeting"
+definition = {
+ "name": name,
+ "nodes": [{
+ "node_instance_id": "hello",
+ "node_id": "n_text_value",
+ "node_version": "1.0.0",
+ "title": "Hello",
+ "params": {"text": "Hello from Python"},
+ }],
+ "edges": [],
+ "graph_outputs": {
+ "text": {"node_instance_id": "hello", "output_name": "text"},
+ },
+}
+
+with Subfork() as client:
+ client.graphs.validate(definition)
+ graph = client.graphs.create(name=name, definition=definition)
+ client.graphs.publish(graph["id"], version="v1")
+
+ execution = client.graphs.execute(graph["id"], version="v1")
+ result = client.executions.wait(execution["id"], timeout=120)
+ print(result["status"], result.get("outputs"))
+```
+
+To test a draft before publishing, call `graphs.execute(graph_id)` without a
+version. To update a draft, use `graphs.update()` with its name, definition, and
+description; omitting the description clears it.
+
+## Reuse building blocks
+
+Explore published graphs and inspect their versioned interfaces:
+
+```python
+with Subfork() as client:
+ blocks = client.graphs.published()
+ block = client.graphs.published_version(graph_id, "v1")
+```
+
+Replace `graph_id` with the ID of a graph you want to use. Published graph versions
+can be composed into larger graphs using the returned composite node manifest.
+Pin child versions so later publications do not change your graph's behavior.
+
+## Handle failures
+
+HTTP failures raise `APIError` subclasses, including `AuthenticationError`,
+`PermissionDeniedError`, `ValidationError`, and `RateLimitError`. These expose
+`status_code` and an optional `retry_after` header.
+
+`executions.wait()` returns failed and canceled runs as well as successful ones,
+so check the returned status. An `ExecutionTimeout` stops polling but does not
+cancel the remote run; use `executions.cancel(execution_id)` to request cancellation.
+
+The client does not automatically retry requests. After a `TransportError`, check
+remote state before repeating a create, publish, or execute request.
+
+## Documentation
+
+See 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.
+
+## Contributing
+
+See [CONTRIBUTING.md](CONTRIBUTING.md) for development and code-quality checks.
+
+## License
+
+[BSD-3-Clause](LICENSE).
diff --git a/assets/subfork-banner.png b/assets/subfork-banner.png
new file mode 100644
index 0000000..2172a9c
Binary files /dev/null and b/assets/subfork-banner.png differ
diff --git a/docs/assets/subfork-banner.png b/docs/assets/subfork-banner.png
new file mode 100644
index 0000000..2172a9c
Binary files /dev/null and b/docs/assets/subfork-banner.png differ
diff --git a/docs/cli.md b/docs/cli.md
new file mode 100644
index 0000000..077144e
--- /dev/null
+++ b/docs/cli.md
@@ -0,0 +1,71 @@
+# Command line
+
+The `subfork` command uses the same `SUBFORK_API_KEY` as the Python client.
+You can also invoke it with `python -m subfork`.
+
+## Find and reuse graphs
+
+```bash
+subfork list
+subfork get GRAPH_ID
+subfork published
+subfork versions GRAPH_ID
+subfork interface GRAPH_ID
+subfork export GRAPH_ID --output graph.json
+```
+
+Export writes the draft definition, ready for `create`. It does not include
+published versions, tags, or the separate graph description.
+
+## Create and publish
+
+```bash
+subfork validate graph.json
+subfork create graph.json --name "My graph"
+subfork publish GRAPH_ID --version v1
+```
+
+Replace `GRAPH_ID` with the ID returned by `create`. Definitions must be JSON
+objects; use `-` instead of a filename to read from stdin. For example files from
+the community repository, see [loading an example](examples.md#load-a-community-example).
+
+## Execute
+
+```bash
+subfork execute GRAPH_ID
+subfork execute GRAPH_ID --version v1 --inputs inputs.json
+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
+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:
+
+```bash
+subfork execute GRAPH_ID > results.json
+```
+
+| Option | Behavior |
+| --- | --- |
+| `--no-wait` | Return the submission ID and status immediately |
+| `--raw` | Return the full execution snapshot |
+| `--wait-timeout 300` | Wait up to 300 seconds; the default is 120 |
+| `--poll-interval 1` | Poll every second; the default is 2 |
+
+A timeout or interruption does not cancel the remote run. Set `NO_COLOR=1` to
+disable colors. Redirected stderr uses plain text without animation.
+
+## Help and exit codes
+
+```bash
+subfork --help
+subfork execute --help
+```
+
+Exit code `0` means success, `1` means an operational error, failed validation,
+unsuccessful execution, or wait timeout, and `2` means invalid command usage.
+Interrupted execution returns `130`.
diff --git a/docs/examples.md b/docs/examples.md
new file mode 100644
index 0000000..e0a81a1
--- /dev/null
+++ b/docs/examples.md
@@ -0,0 +1,71 @@
+# Examples
+
+Explore [examples.subfork.com](https://examples.subfork.com) for graphs covering
+AI, data pipelines, dashboards, media, and more. Download definitions and read
+their notes in the [subfork-examples repository](https://github.com/subforkdev/subfork-examples).
+It's also the place to contribute examples for other builders.
+
+## Hello from Python
+
+This graph returns a string. It requires read, write, and run scopes.
+
+```python
+from subfork import Subfork
+
+name = "Python greeting"
+definition = {
+ "name": name,
+ "nodes": [{
+ "node_instance_id": "hello",
+ "node_id": "n_text_value",
+ "node_version": "1.0.0",
+ "title": "Hello",
+ "params": {"text": "Hello from Python"},
+ }],
+ "edges": [],
+ "graph_outputs": {
+ "text": {"node_instance_id": "hello", "output_name": "text"},
+ },
+}
+
+with Subfork() as client:
+ graph = client.graphs.create(name=name, definition=definition)
+ execution = client.graphs.execute(graph["id"])
+ result = client.executions.wait(execution["id"])
+ print(result["status"], result.get("outputs"))
+```
+
+This creates a public graph and uses execution quota. To publish a reusable
+version, call `client.graphs.publish(graph["id"], version="v1")` with publish scope.
+
+## Load a community example
+
+Download a `*.subfork.json` file from the examples repository and read its
+adjacent `*.notes.md` file for required inputs, secrets, and provider setup.
+These exports wrap the graph definition in a `definition` field:
+
+```python
+import json
+from subfork import Subfork
+
+with open("example.subfork.json", encoding="utf-8") as stream:
+ exported = json.load(stream)
+definition = exported["definition"]
+
+with Subfork() as client:
+ validation = client.graphs.validate(definition)
+ if not validation.get("valid", False):
+ raise ValueError(validation)
+ graph = client.graphs.create(
+ name=definition.get("name", "Imported example"),
+ description=definition.get("description", ""),
+ definition=definition,
+ )
+ print(graph["id"])
+```
+
+Configure any required secrets in the graph's settings before executing it.
+The CLI accepts a raw definition, so extract `exported["definition"]` to a JSON
+file before using `subfork create`. CLI exports already use that raw format.
+
+Browse the [example catalog](https://examples.subfork.com) to choose your next graph.
diff --git a/docs/index.md b/docs/index.md
new file mode 100644
index 0000000..f9eed20
--- /dev/null
+++ b/docs/index.md
@@ -0,0 +1,25 @@
+
+
+# Subfork Python
+
+Create, publish, and run [Subfork](https://subfork.com) graphs from Python or your terminal.
+
+```bash
+pip install subfork
+```
+
+Start with [installation and authentication](installation.md), then follow the
+[Python guide](usage.md) or [command-line guide](cli.md).
+
+## Find your next graph
+
+Browse [Subfork Examples](https://examples.subfork.com) for ready-made graphs covering
+data pipelines, dashboards, AI, media, and more. The
+[subfork-examples repository](https://github.com/subforkdev/subfork-examples) contains
+the graph files and notes you can learn from, adapt, and contribute to.
+
+Our [examples guide](examples.md) shows how to run your first graph and load a
+community example with Python.
+
+[Source code](https://github.com/subforkdev/subfork-python) ·
+[Report an issue](https://github.com/subforkdev/subfork-python/issues)
diff --git a/docs/installation.md b/docs/installation.md
new file mode 100644
index 0000000..c8da813
--- /dev/null
+++ b/docs/installation.md
@@ -0,0 +1,47 @@
+# Installation
+
+Requires Python 3.8 or later.
+
+```bash
+pip install subfork
+```
+
+This installs the Python package and the `subfork` command. Both connect to
+[subfork.com](https://subfork.com).
+
+## Set your API key
+
+Create a key in **Account → API keys**, then set it in your environment:
+
+```bash
+export SUBFORK_API_KEY="your-api-key"
+subfork list
+```
+
+The client reads exported environment variables; it does not load `.env` files
+automatically. Keep your key out of source files and graph definitions.
+
+## Choose permissions
+
+| Scope | Allows |
+| --- | --- |
+| `graphs:read` | Read accessible graphs, discover nodes, and inspect your executions |
+| `graphs:write` | Create and update your graphs, and validate definitions |
+| `graphs:run` | Execute your graphs and cancel your executions |
+| `graphs:publish` | Publish versions of your graphs |
+
+The CLI's default execution workflow needs both read and run permissions because
+it waits for the result. Keys act as their owning account; scopes do not grant
+permission to modify other users' graphs. Public graphs can be read and reused.
+
+## Check the connection
+
+```python
+from subfork import Subfork
+
+with Subfork() as client:
+ for graph in client.graphs.list():
+ print(graph["id"], graph["name"])
+```
+
+Continue with [Python usage](usage.md) or the [CLI](cli.md).
diff --git a/docs/mkpages.yml b/docs/mkpages.yml
new file mode 100644
index 0000000..1e70b59
--- /dev/null
+++ b/docs/mkpages.yml
@@ -0,0 +1,20 @@
+title: Subfork Python
+description: Create, publish, and run Subfork graphs from Python and the command line.
+theme: dark
+navigation:
+ - label: Home
+ href: /
+ - label: Install
+ href: /installation/
+ - label: Python
+ href: /usage/
+ - label: CLI
+ href: /cli/
+ - label: Examples
+ href: /examples/
+ - label: Troubleshooting
+ href: /troubleshooting/
+ - label: Example Catalog
+ href: https://examples.subfork.com
+ - label: GitHub
+ href: https://github.com/subforkdev/subfork-python
diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md
new file mode 100644
index 0000000..093730b
--- /dev/null
+++ b/docs/troubleshooting.md
@@ -0,0 +1,41 @@
+# Troubleshooting
+
+## 401: authentication failed
+
+Check that `SUBFORK_API_KEY` is exported and is an active key from your Subfork
+account. The client does not automatically read `.env` files. Create a replacement
+under **Account → API keys** if the key has expired or been revoked.
+
+## 403: operation not permitted
+
+Check the key's scopes. Waiting for execution results needs `graphs:read` as well
+as `graphs:run`. See [permissions](installation.md#choose-permissions).
+
+## 404: graph not found
+
+Check the graph ID and ownership. Reading a public graph does not give you
+permission to modify, publish, or execute the original. Create your own copy first.
+
+## 409: graph secret cannot be decrypted
+
+Re-enter and save the affected secret in **Graph Settings → Secrets**, then retry.
+This is a stored graph credential, separate from the key used by the Python
+client. Other 409 responses can indicate a different execution setup conflict.
+
+## Output file already exists
+
+Use `-f` / `--force` to overwrite it:
+
+```bash
+subfork execute GRAPH_ID -o results.json --force
+```
+
+## Execution timed out or failed
+
+A wait timeout stops the client from polling; the run may still be active. Use
+`client.executions.get(execution_id)` to inspect it, or open the graph in Subfork.
+Use `client.executions.cancel(execution_id)` if you want to cancel it.
+
+A terminal `failed` state is different from a timeout. Inspect the execution
+snapshot or the graph's execution history for node errors. The CLI's `--raw`
+option includes the full snapshot when submitting a run.
diff --git a/docs/usage.md b/docs/usage.md
new file mode 100644
index 0000000..c62bdef
--- /dev/null
+++ b/docs/usage.md
@@ -0,0 +1,89 @@
+# Python usage
+
+Use `Subfork` as a context manager to close its HTTP connection when finished.
+Responses are ordinary dictionaries and lists.
+
+```python
+from subfork import Subfork
+
+with Subfork() as client:
+ graphs = client.graphs.list()
+ nodes = client.nodes.list()
+```
+
+## Create and publish
+
+Load a graph definition from JSON, validate it, and save a new graph:
+
+```python
+import json
+from subfork import Subfork
+
+with open("graph.json", encoding="utf-8") as stream:
+ definition = json.load(stream)
+
+with Subfork() as client:
+ validation = client.graphs.validate(definition)
+ if not validation.get("valid", False):
+ raise ValueError(validation)
+ graph = client.graphs.create(name="My graph", definition=definition)
+ published = client.graphs.publish(graph["id"], version="v1")
+```
+
+Graphs are public. Publication creates an immutable version. To replace draft
+content, use `client.graphs.update(graph_id, name=..., definition=..., description=...)`.
+Supply the description you want to retain; omitting it clears the description.
+
+## Execute and wait
+
+```python
+with Subfork() as client:
+ execution = client.graphs.execute(graph_id, version="v1", inputs={})
+ result = client.executions.wait(execution["id"], timeout=120)
+ if result["status"] == "completed":
+ print(result["outputs"])
+ else:
+ print("Execution ended:", result["status"])
+```
+
+Replace `graph_id` with your graph ID and `inputs` with its named inputs. Omit
+`version` to run the draft. Runs consume your account's execution quota.
+
+`execute()` submits immediately; `wait()` polls until a terminal state. You can
+also call `executions.get(execution_id)` or `executions.cancel(execution_id)`.
+A wait timeout stops polling without canceling the remote run.
+
+## Discover reusable graphs
+
+```python
+with Subfork() as client:
+ catalog = client.graphs.published()
+ block = client.graphs.published_version(graph_id, "v1")
+ interface = client.graphs.interface(owned_graph_id)
+```
+
+Published graph versions include a composite node manifest for use in larger
+graphs. Pin versions to keep later publications from changing your graph.
+See [examples](examples.md) for complete definitions.
+
+## Handle errors
+
+```python
+from subfork import APIError, ExecutionTimeout, Subfork
+
+try:
+ with Subfork() as client:
+ execution = client.graphs.execute(graph_id)
+ result = client.executions.wait(execution["id"])
+except APIError as error:
+ print(error.status_code, str(error))
+except ExecutionTimeout as error:
+ print("Still waiting for:", error.execution_id)
+```
+
+API errors include `AuthenticationError`, `PermissionDeniedError`, `NotFoundError`,
+`ValidationError`, and `RateLimitError`. The client does not automatically retry
+requests. After a `TransportError`, check remote state before submitting another
+create, publish, or execute request.
+
+See [troubleshooting](troubleshooting.md) for common errors.
diff --git a/pyproject.toml b/pyproject.toml
new file mode 100644
index 0000000..ce5da3a
--- /dev/null
+++ b/pyproject.toml
@@ -0,0 +1,76 @@
+[build-system]
+requires = ["setuptools>=61", "wheel"]
+build-backend = "setuptools.build_meta"
+
+[project]
+name = "subfork"
+version = "2.0.0"
+description = "Python API client for Subfork graph authoring and execution"
+readme = "README.md"
+license = {file = "LICENSE"}
+requires-python = ">=3.8"
+dependencies = ["httpx>=0.28,<1"]
+classifiers = [
+ "Development Status :: 3 - Alpha",
+ "Intended Audience :: Developers",
+ "Programming Language :: Python :: 3",
+ "Typing :: Typed",
+]
+
+[project.urls]
+Documentation = "https://python.subfork.com"
+
+[project.scripts]
+subfork = "subfork.cli:main"
+
+[project.optional-dependencies]
+docs = ["mkpages==0.4.1"]
+dev = [
+ "pytest>=8,<10",
+ "build>=1,<2",
+ "twine>=6,<7",
+ "flake8==7.1.1",
+ "flake8-docstrings==1.7.0",
+ "mccabe==0.7.0",
+ "isort==5.13.2",
+ "black==24.8.0",
+ "mypy==1.14.1",
+]
+
+[tool.setuptools]
+license-files = ["LICENSE"]
+
+[tool.setuptools.packages.find]
+where = ["src"]
+
+[tool.setuptools.package-data]
+subfork = ["py.typed"]
+
+[tool.pytest.ini_options]
+testpaths = ["tests"]
+
+[tool.isort]
+profile = "black"
+line_length = 100
+skip_gitignore = true
+
+[tool.black]
+line-length = 100
+target-version = ["py38"]
+extend-exclude = '/\.venv[^/]*/'
+
+[tool.mypy]
+python_version = "3.8"
+files = ["src/subfork"]
+disallow_untyped_defs = true
+check_untyped_defs = true
+no_implicit_optional = true
+warn_unused_ignores = true
+warn_redundant_casts = true
+show_error_codes = true
+
+# Newer environments install CLI/async dependency internals with newer syntax.
+# Skip those internals while retaining HTTPX's client types and our 3.8 target.
+[[tool.mypy.overrides]]
+module = ["httpx._main", "anyio", "anyio.*"]
+follow_imports = "skip"
diff --git a/scripts/bump-version.py b/scripts/bump-version.py
new file mode 100644
index 0000000..c618130
--- /dev/null
+++ b/scripts/bump-version.py
@@ -0,0 +1,76 @@
+"""Keep package metadata and the public client version synchronized."""
+
+import argparse
+import os
+import re
+from pathlib import Path
+from typing import Dict, Optional, Sequence, Tuple
+
+ROOT = Path(__file__).resolve().parents[1]
+SEMVER = re.compile(r"(0|[1-9][0-9]*)\.(0|[1-9][0-9]*)\.(0|[1-9][0-9]*)")
+
+
+def version_files(root: Path) -> Tuple[str, Dict[Path, Tuple[str, re.Match]]]:
+ """Read both declarations and reject missing, invalid, or mismatched versions."""
+ files = {}
+ versions = []
+ for relative, pattern in (
+ ("pyproject.toml", r'(?ms)^\[project\]\s*\n(?:(?!^\[).)*?^version = "([^"]+)"'),
+ ("src/subfork/__init__.py", r'(?m)^__version__ = "([^"]+)"'),
+ ):
+ path = root / relative
+ content = path.read_text(encoding="utf-8")
+ matches = list(re.finditer(pattern, content))
+ if len(matches) != 1 or not SEMVER.fullmatch(matches[0].group(1)):
+ raise ValueError("Expected one stable x.y.z version in {}.".format(relative))
+ files[path] = (content, matches[0])
+ versions.append(matches[0].group(1))
+ if versions[0] != versions[1]:
+ raise ValueError("Version drift between pyproject.toml and __version__; no files changed.")
+ return versions[0], files
+
+
+def next_version(current: str, requested: str) -> str:
+ """Compute a stable release version, resetting lower components on major/minor bumps."""
+ if SEMVER.fullmatch(requested):
+ return requested
+ major, minor, patch = map(int, current.split("."))
+ if requested == "major":
+ return "{}.0.0".format(major + 1)
+ if requested == "minor":
+ return "{}.{}.0".format(major, minor + 1)
+ if requested == "patch":
+ return "{}.{}.{}".format(major, minor, patch + 1)
+ raise ValueError("Use major, minor, patch, or a stable x.y.z version.")
+
+
+def main(argv: Optional[Sequence[str]] = None) -> int:
+ """Show, check, or update the version without committing or tagging a release."""
+ parser = argparse.ArgumentParser(description=__doc__)
+ group = parser.add_mutually_exclusive_group(required=True)
+ group.add_argument("--show", action="store_true")
+ group.add_argument("--check", action="store_true")
+ group.add_argument("--bump", choices=("major", "minor", "patch"))
+ group.add_argument("--set", action="store_true", help="Read the explicit version from VERSION")
+ args = parser.parse_args(argv)
+ try:
+ current, files = version_files(ROOT)
+ if args.show or args.check:
+ print(current if args.show else "Version fields match: " + current)
+ return 0
+ requested = os.environ.get("VERSION", "") if args.set else args.bump
+ if not requested:
+ raise ValueError("Usage: make bump-version VERSION=x.y.z")
+ version = next_version(current, requested)
+ for path, (content, match) in files.items():
+ updated = content[: match.start(1)] + version + content[match.end(1) :]
+ if updated != content:
+ path.write_text(updated, encoding="utf-8")
+ print("{} -> {}".format(current, version))
+ return 0
+ except (OSError, ValueError) as error:
+ parser.exit(1, str(error) + "\n")
+
+
+if __name__ == "__main__":
+ raise SystemExit(main())
diff --git a/src/subfork/__init__.py b/src/subfork/__init__.py
new file mode 100644
index 0000000..b59ac78
--- /dev/null
+++ b/src/subfork/__init__.py
@@ -0,0 +1,30 @@
+"""Subfork's Python API client. BYO worker support is a separate distribution."""
+
+from .client import Subfork
+from .errors import (
+ APIError,
+ AuthenticationError,
+ ExecutionTimeout,
+ InvalidResponseError,
+ NotFoundError,
+ PermissionDeniedError,
+ RateLimitError,
+ SubforkError,
+ TransportError,
+ ValidationError,
+)
+
+__version__ = "2.0.0"
+__all__ = [
+ "Subfork",
+ "SubforkError",
+ "APIError",
+ "AuthenticationError",
+ "PermissionDeniedError",
+ "NotFoundError",
+ "ValidationError",
+ "RateLimitError",
+ "TransportError",
+ "InvalidResponseError",
+ "ExecutionTimeout",
+]
diff --git a/src/subfork/__main__.py b/src/subfork/__main__.py
new file mode 100644
index 0000000..ffa012f
--- /dev/null
+++ b/src/subfork/__main__.py
@@ -0,0 +1,6 @@
+"""Support invocation with python -m subfork."""
+
+from .cli import main
+
+if __name__ == "__main__":
+ raise SystemExit(main())
diff --git a/src/subfork/cli.py b/src/subfork/cli.py
new file mode 100644
index 0000000..e982bc0
--- /dev/null
+++ b/src/subfork/cli.py
@@ -0,0 +1,386 @@
+"""Command-line graph authoring and execution with JSON input and output."""
+
+import argparse
+import json
+import math
+import os
+import shutil
+import sys
+import threading
+from contextlib import contextmanager
+from pathlib import Path
+from typing import Any, Callable, Dict, Iterator, Optional, Sequence
+
+from . import APIError, AuthenticationError, ExecutionTimeout, Subfork, SubforkError, __version__
+
+
+def print_api_error(status_code: int, message: str) -> None:
+ """Write a concise API diagnostic, coloring only the status code on terminals."""
+ code = str(status_code)
+ if sys.stderr.isatty() and os.environ.get("TERM") != "dumb" and not os.environ.get("NO_COLOR"):
+ code = "\033[33m" + code + "\033[0m"
+ print("{}: {}".format(code, message), file=sys.stderr)
+
+
+@contextmanager
+def execution_status(graph_name: str) -> Iterator[Callable[[Dict[str, Any]], None]]:
+ """Animate the active nodes on terminal stderr while preserving JSON stdout."""
+ stream = sys.stderr
+ terminal = stream.isatty() and os.environ.get("TERM") != "dumb"
+ lock = threading.Lock()
+ label = "Graph " + graph_name
+ status = "Running"
+ titles: Dict[str, str] = {}
+ reported: Dict[str, str] = {}
+ name_width = len(label)
+ stopped = threading.Event()
+ completed_marker = "✓"
+ try:
+ completed_marker.encode(stream.encoding or "ascii")
+ except UnicodeError:
+ completed_marker = "+"
+ frames = "⠋⠙⠹⠸⠼⠴⠦⠧⠇⠏"
+ try:
+ frames.encode(stream.encoding or "ascii")
+ except UnicodeError:
+ frames = "|/-\\"
+
+ def update(snapshot: Dict[str, Any]) -> None:
+ """Replace the display state with active node titles from a snapshot."""
+ nonlocal label, status, name_width
+ definition = snapshot.get("definition") or {}
+ for node in definition.get("nodes", []):
+ titles[node["node_instance_id"]] = node.get("title") or node["node_instance_id"]
+ nodes = snapshot.get("node_executions") or {}
+ active = [
+ titles.get(key, key) for key, node in nodes.items() if node.get("status") == "running"
+ ]
+ with lock:
+ name_width = max(
+ [name_width]
+ + [len("Node " + title) for title in titles.values()]
+ + [len("Node " + key) for key in nodes if key not in titles]
+ )
+ for key, node in nodes.items():
+ outcome = node.get("status")
+ if outcome not in {
+ "completed",
+ "failed",
+ "canceled",
+ "cancelled",
+ "skipped",
+ "outcome_unknown",
+ }:
+ continue
+ if reported.get(key) == outcome:
+ continue
+ reported[key] = outcome
+ state = "Canceled" if outcome == "cancelled" else outcome.replace("_", " ").title()
+ name = "Node " + titles.get(key, key)
+ if terminal:
+ stream.write(
+ "\r\033[2K"
+ + format_line(
+ completed_marker if outcome == "completed" else "-", name, state
+ )
+ + "\n"
+ )
+ else:
+ stream.write(
+ format_line(
+ completed_marker if outcome == "completed" else "-", name, state
+ )
+ + "\n"
+ )
+ stream.flush()
+ label = (
+ ("Node " if len(active) == 1 else "Nodes ") + ", ".join(active)
+ if active
+ else "Graph " + graph_name
+ )
+ status = (
+ "Running"
+ if active
+ else str(snapshot.get("status", "running")).replace("_", " ").title()
+ )
+
+ def format_line(frame: str, name: str, state: str) -> str:
+ """Format a width-limited progress row without splitting color escapes."""
+ name = "".join(char if char.isprintable() else " " for char in name)
+ state = "".join(char if char.isprintable() else " " for char in state)
+ width = max(1, shutil.get_terminal_size().columns - 1)
+ # Reserve space for the longest status so different outcomes align too.
+ status_column = min(name_width + 14, max(7, width - len("Outcome Unknown")))
+ name = name[: max(0, status_column - 7)]
+ dots = "." * max(3, status_column - len(name) - 4)
+ plain = "{} {} {} {}".format(frame, name, dots, state)
+ if len(plain) > width:
+ return plain[:width]
+ if not terminal or os.environ.get("NO_COLOR"):
+ return plain
+ color = "32" if state in {"Running", "Completed"} else "31" if state == "Failed" else "33"
+ colored_state = "\033[" + color + "m" + state + "\033[0m"
+ marker_color = "32" if state == "Completed" else "33"
+ return "\033[{}m{}\033[0m {} {} {}".format(marker_color, frame, name, dots, colored_state)
+
+ def render(index: int) -> None:
+ """Draw the active row without interleaving retained completion lines."""
+ with lock:
+ stream.write("\r\033[2K" + format_line(frames[index % len(frames)], label, status))
+ stream.flush()
+
+ def animate() -> None:
+ """Refresh the spinner independently of network requests and polling."""
+ index = 1
+ while not stopped.wait(0.1):
+ try:
+ render(index)
+ except (OSError, ValueError):
+ return
+ index += 1
+
+ if not terminal:
+ safe_name = "".join(char if char.isprintable() else " " for char in graph_name)
+ print("Graph {} .......... Running".format(safe_name), file=stream)
+ yield update
+ return
+ worker = threading.Thread(target=animate, name="subfork-spinner", daemon=True)
+ try:
+ render(0)
+ worker.start()
+ yield update
+ finally:
+ stopped.set()
+ if worker.ident is not None:
+ worker.join()
+ stream.write("\r\033[2K")
+ stream.flush()
+
+
+def build_parser() -> argparse.ArgumentParser:
+ """Build the public command tree without opening a connection."""
+ parser = argparse.ArgumentParser(description=__doc__)
+ parser.add_argument("--version", action="version", version=__version__)
+ parser.add_argument("--base-url", default=os.getenv("SUBFORK_BASE_URL", "https://subfork.com"))
+ parser.add_argument("--timeout", type=float, default=30, help="HTTP timeout in seconds")
+ commands = parser.add_subparsers(dest="command", required=True)
+ commands.add_parser("list", help="List account graphs")
+ commands.add_parser("published", help="List published building blocks")
+ for command in ("get", "versions", "interface", "export", "publish", "execute"):
+ child = commands.add_parser(command)
+ child.add_argument("graph_id")
+ if command in {"export", "execute"}:
+ child.add_argument(
+ "-f", "--force", action="store_true", help="Overwrite an existing output file"
+ )
+ if command == "export":
+ child.add_argument(
+ "--output", default="-", help="Definition JSON path, or - for stdout"
+ )
+ elif command == "publish":
+ child.add_argument("--version", required=True, help="New immutable version, e.g. v1")
+ child.add_argument("--interface", help="Optional interface JSON file")
+ 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(
+ "--no-wait", action="store_true", help="Return submission status immediately"
+ )
+ child.add_argument(
+ "--raw", action="store_true", help="Print the full execution snapshot"
+ )
+ child.add_argument(
+ "--wait-timeout",
+ type=float,
+ default=120,
+ help="Maximum wait in seconds (default: 120)",
+ )
+ child.add_argument(
+ "--poll-interval",
+ type=float,
+ default=2,
+ help="Seconds between status checks (default: 2)",
+ )
+ child.add_argument("--inputs", help="Inputs JSON file, or - for stdin")
+ for command in ("validate", "create"):
+ child = commands.add_parser(command)
+ child.add_argument("file", help="Graph definition JSON file, or - for stdin")
+ child.add_argument("--name", help="Override the definition's name")
+ if command == "create":
+ child.add_argument("--description", default="")
+ child.add_argument("--tag", action="append", default=[])
+ return parser
+
+
+def read_object(filename: str) -> Dict[str, Any]:
+ """Read a JSON object from a UTF-8 file or standard input."""
+ try:
+ content = (
+ sys.stdin.read() if filename == "-" else Path(filename).read_text(encoding="utf-8")
+ )
+ value = json.loads(content)
+ except (ValueError, UnicodeError) as exc:
+ raise ValueError("Input must contain valid UTF-8 JSON.") from exc
+ if not isinstance(value, dict):
+ raise ValueError("Input JSON must be an object.")
+ return value
+
+
+def graph_command(client: Subfork, args: argparse.Namespace) -> Any:
+ """Dispatch graph operations through the same SDK used by Python callers."""
+ graphs = client.graphs
+ if args.command == "list":
+ return graphs.list()
+ if args.command == "published":
+ return graphs.published()
+ if args.command in {"create", "validate"}:
+ definition = read_object(args.file)
+ name = args.name if args.name is not None else definition.get("name", "Untitled Graph")
+ if not isinstance(name, str) or not name.strip():
+ raise ValueError("Graph name must be a nonempty string.")
+ if args.command == "validate":
+ return graphs.validate(definition, name=name)
+ return graphs.create(
+ name=name, definition=definition, description=args.description, tags=args.tag
+ )
+ if args.command == "get":
+ return graphs.get(args.graph_id)
+ if args.command == "export":
+ graph = graphs.get(args.graph_id)
+ exported = graph.get("definition")
+ if not isinstance(exported, dict):
+ raise ValueError("Graph response does not contain a definition object.")
+ return exported
+ if args.command == "versions":
+ return graphs.versions(args.graph_id)
+ if args.command == "interface":
+ return graphs.interface(args.graph_id)
+ if args.command == "publish":
+ interface = read_object(args.interface) if args.interface else None
+ return graphs.publish(
+ args.graph_id, version=args.version, interface=interface, comment=args.comment
+ )
+ inputs = read_object(args.inputs) if args.inputs else None
+ if any(
+ not math.isfinite(value) or value <= 0 for value in (args.wait_timeout, args.poll_interval)
+ ):
+ raise ValueError("Wait timeout and poll interval must be positive finite numbers.")
+ result = graphs.execute(args.graph_id, version=args.version, inputs=inputs)
+ terminal_statuses = {
+ "completed",
+ "failed",
+ "canceled",
+ "cancelled",
+ "outcome_unknown",
+ }
+ if not args.no_wait:
+ execution_id = result.get("id") or result.get("execution_id")
+ if not isinstance(execution_id, str) or not execution_id:
+ raise ValueError(
+ "Submission returned no execution ID; inspect remote state before retrying."
+ )
+ snapshot_definition = result.get("definition")
+ graph_name = (
+ snapshot_definition.get("name") if isinstance(snapshot_definition, dict) else None
+ )
+ with execution_status(
+ graph_name if isinstance(graph_name, str) else args.graph_id
+ ) as update:
+ update(result)
+ if result.get("status") not in terminal_statuses:
+ result = client.executions.wait(
+ execution_id,
+ timeout=args.wait_timeout,
+ poll_interval=args.poll_interval,
+ on_update=update,
+ )
+ return result
+
+
+def main(argv: Optional[Sequence[str]] = None) -> int:
+ """Run the CLI; return zero on success and one on operational failure.
+
+ Argparse reports usage errors with exit code two. Failed validation returns
+ one, retaining its JSON diagnostics on stdout.
+ """
+ args = build_parser().parse_args(argv)
+ try:
+ result_file = getattr(args, "out", None)
+ force = getattr(args, "force", False)
+ export_file = getattr(args, "output", "-")
+ destination = (
+ result_file
+ 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:
+ 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)
+ failed = args.command == "execute" and result.get("status") in {
+ "failed",
+ "canceled",
+ "cancelled",
+ "outcome_unknown",
+ }
+ display = result
+ if args.command == "execute" and not args.raw:
+ if args.no_wait or failed:
+ display = {
+ "execution_id": result.get("id") or result.get("execution_id"),
+ "status": result.get("status"),
+ }
+ 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)
+ if failed:
+ print(
+ "subfork: execution did not complete successfully; use --raw for details",
+ file=sys.stderr,
+ )
+ return 1
+ if args.command == "validate" and result.get("valid") is False:
+ return 1
+ return 0
+ except ExecutionTimeout as exc:
+ print(
+ "subfork: execution {} timed out while waiting; the remote run is not canceled".format(
+ exc.execution_id
+ ),
+ file=sys.stderr,
+ )
+ return 1
+ except AuthenticationError:
+ print_api_error(
+ 401,
+ "Authentication failed. Check that SUBFORK_API_KEY "
+ "is an active key issued by the service selected with --base-url or "
+ "SUBFORK_BASE_URL (default: https://subfork.com). "
+ "The CLI reads exported environment variables; it does not load .env files.",
+ )
+ return 1
+ except APIError as exc:
+ message = str(exc)
+ prefix = "Subfork API returned HTTP {}.".format(exc.status_code)
+ if message.startswith(prefix):
+ message = message[len(prefix) :].strip() or "API request failed."
+ print_api_error(exc.status_code, message)
+ return 1
+ except (SubforkError, ValueError, OSError) as exc:
+ print("subfork: {}".format(exc), file=sys.stderr)
+ return 1
+ except KeyboardInterrupt:
+ print(
+ "subfork: interrupted; remote execution is not canceled automatically", file=sys.stderr
+ )
+ return 130
diff --git a/src/subfork/client.py b/src/subfork/client.py
new file mode 100644
index 0000000..d3209a3
--- /dev/null
+++ b/src/subfork/client.py
@@ -0,0 +1,163 @@
+"""Synchronous, explicitly closed HTTP client. No automatic write retries."""
+
+from __future__ import annotations
+
+import math
+import os
+import re
+from types import TracebackType
+from typing import Any, Optional, Type
+
+import httpx
+
+from .errors import (
+ APIError,
+ AuthenticationError,
+ InvalidResponseError,
+ NotFoundError,
+ PermissionDeniedError,
+ RateLimitError,
+ TransportError,
+ ValidationError,
+)
+from .resources import Executions, Graphs, Nodes
+
+
+def _error_message(response: httpx.Response) -> str:
+ """Translate recognized server errors into fixed, credential-safe guidance."""
+ message = "Subfork API returned HTTP %s." % response.status_code
+ if response.status_code != 409:
+ return message
+ try:
+ payload = response.json()
+ except ValueError:
+ return message
+ detail = payload.get("detail") if isinstance(payload, dict) else None
+ if isinstance(detail, str) and re.fullmatch(
+ r"Graph secret '[^\r\n]*' could not be decrypted\. The server encryption "
+ r"key may have changed; re-enter this secret in Graph Settings\.",
+ detail,
+ ):
+ return (
+ message + " A stored graph secret could not be decrypted. "
+ "The server encryption key may have changed. "
+ "Re-enter and save the affected secret in Graph Settings > Secrets, then retry."
+ )
+ return message
+
+
+class Subfork:
+ """Manage authenticated HTTP requests and graph resources.
+
+ The client owns its connection pool. Use a context manager or close() to
+ release it. Writes are never automatically retried.
+ """
+
+ def __init__(
+ self,
+ api_key: Optional[str] = None,
+ *,
+ base_url: str = "https://subfork.com",
+ timeout: float = 30,
+ transport: Optional[httpx.BaseTransport] = None,
+ ) -> None:
+ """Create an authenticated client for the Subfork API.
+
+ Args:
+ api_key: Bearer secret, or None to read SUBFORK_API_KEY.
+ base_url: API origin, defaulting to https://subfork.com, without
+ an API path, query, or credentials.
+ timeout: Positive per-I/O-phase HTTP timeout in seconds.
+ transport: Optional HTTPX transport, primarily for testing.
+
+ Raises:
+ ValueError: Credentials, origin, or timeout are invalid.
+ """
+ key = os.environ.get("SUBFORK_API_KEY") if api_key is None else api_key
+ if not key or any(char.isspace() for char in key):
+ raise ValueError(
+ "Provide an API key or set SUBFORK_API_KEY; keys cannot contain whitespace."
+ )
+ if not math.isfinite(timeout) or timeout <= 0:
+ raise ValueError("timeout must be a positive finite number.")
+ origin = httpx.URL(base_url)
+ if (
+ not origin.host
+ or origin.userinfo
+ or origin.query
+ or origin.fragment
+ or origin.path not in {"", "/"}
+ ):
+ raise ValueError(
+ "base_url must be an origin without credentials, path, query or fragment."
+ )
+ if origin.scheme != "https" and not (
+ origin.scheme == "http"
+ and (
+ origin.host in {"localhost", "127.0.0.1", "::1"}
+ or origin.host.endswith(".localhost")
+ )
+ ):
+ raise ValueError("Use HTTPS; HTTP is allowed only for localhost.")
+ self.timeout = timeout
+ self._client = httpx.Client(
+ base_url=str(origin).rstrip("/") + "/api/v1/",
+ timeout=timeout,
+ headers={"Authorization": "Bearer " + key, "Accept": "application/json"},
+ follow_redirects=False,
+ transport=transport,
+ )
+ self.nodes = Nodes(self)
+ self.graphs = Graphs(self)
+ self.executions = Executions(self)
+
+ def _request(self, method: str, path: str, **kwargs: Any) -> Any:
+ """Send one request and decode JSON, returning None for HTTP 204.
+
+ Raise APIError subclasses for HTTP failures, TransportError for network
+ failures, and InvalidResponseError for invalid JSON. Error messages
+ never include response bodies or credentials.
+ """
+ try:
+ response = self._client.request(method, path.lstrip("/"), **kwargs)
+ except httpx.TransportError:
+ raise TransportError(
+ "API request could not be completed; inspect remote state before retrying a write."
+ ) from None
+ if not response.is_success:
+ errors = {
+ 401: AuthenticationError,
+ 403: PermissionDeniedError,
+ 404: NotFoundError,
+ 422: ValidationError,
+ 429: RateLimitError,
+ }
+ # Fixed messages avoid reflecting secrets from an untrusted response.
+ raise errors.get(response.status_code, APIError)(
+ _error_message(response),
+ status_code=response.status_code,
+ retry_after=response.headers.get("retry-after"),
+ )
+ if response.status_code == 204:
+ return None
+ try:
+ return response.json()
+ except ValueError:
+ raise InvalidResponseError("API returned invalid JSON.") from None
+
+ def close(self) -> None:
+ """Release the HTTP connection pool."""
+ self._client.close()
+
+ def __enter__(self) -> Subfork:
+ """Return this client for use in a with statement."""
+ return self
+
+ def __exit__(
+ self,
+ exc_type: Optional[Type[BaseException]],
+ exc_value: Optional[BaseException],
+ traceback: Optional[TracebackType],
+ ) -> None:
+ """Close the client without suppressing exceptions."""
+ self.close()
diff --git a/src/subfork/errors.py b/src/subfork/errors.py
new file mode 100644
index 0000000..c446f15
--- /dev/null
+++ b/src/subfork/errors.py
@@ -0,0 +1,68 @@
+"""Public client errors; credentials and raw responses are never attached."""
+
+from typing import Optional
+
+
+class SubforkError(Exception):
+ """Base class for client errors."""
+
+
+class TransportError(SubforkError):
+ """The request could not be completed; a write's outcome may be unknown."""
+
+
+class APIError(SubforkError):
+ """Represent an HTTP failure without retaining a raw request or response."""
+
+ def __init__(
+ self, message: str, *, status_code: int, retry_after: Optional[str] = None
+ ) -> None:
+ """Store a safe message, HTTP status, and optional Retry-After header."""
+ super().__init__(message)
+ self.status_code = status_code
+ self.retry_after = retry_after
+
+
+class AuthenticationError(APIError):
+ """Report HTTP 401 for missing, expired, or invalid credentials."""
+
+ pass
+
+
+class PermissionDeniedError(APIError):
+ """Report HTTP 403 when the credential cannot perform an operation."""
+
+ pass
+
+
+class NotFoundError(APIError):
+ """Report HTTP 404 for a missing or inaccessible resource."""
+
+ pass
+
+
+class ValidationError(APIError):
+ """Report HTTP 422 when server request validation fails."""
+
+ pass
+
+
+class RateLimitError(APIError):
+ """Report HTTP 429, with an optional Retry-After header."""
+
+ pass
+
+
+class InvalidResponseError(SubforkError):
+ """Report a successful HTTP response containing invalid JSON."""
+
+ pass
+
+
+class ExecutionTimeout(SubforkError):
+ """Polling expired. The remote execution was not cancelled."""
+
+ def __init__(self, execution_id: str) -> None:
+ """Identify the execution whose polling deadline expired."""
+ super().__init__("Execution polling timed out; the remote execution may still be running.")
+ self.execution_id = execution_id
diff --git a/src/subfork/py.typed b/src/subfork/py.typed
new file mode 100644
index 0000000..e69de29
diff --git a/src/subfork/resources.py b/src/subfork/resources.py
new file mode 100644
index 0000000..4cd1349
--- /dev/null
+++ b/src/subfork/resources.py
@@ -0,0 +1,228 @@
+"""Small resource wrappers over the existing versioned HTTP API."""
+
+from __future__ import annotations
+
+import math
+import time
+from typing import TYPE_CHECKING, Any, Callable, Dict, List, Optional
+from urllib.parse import quote
+
+from .errors import ExecutionTimeout
+
+if TYPE_CHECKING:
+ from .client import Subfork
+
+
+def segment(value: str) -> str:
+ """Encode an identifier as a path segment, rejecting empty or dot segments."""
+ if not value or value in {".", ".."}:
+ raise ValueError("Resource identifiers must be non-empty and cannot be dot segments.")
+ return quote(value, safe="")
+
+
+class Resource:
+ """Share a client connection with a resource-specific wrapper."""
+
+ def __init__(self, client: Subfork) -> None:
+ """Bind this wrapper to its owning client."""
+ self._client = client
+
+
+class Nodes(Resource):
+ """Discover available node implementations and manifests."""
+
+ def list(self, *, include_all_versions: bool = False) -> List[Dict[str, Any]]:
+ """Return node manifests, optionally including older published versions."""
+ return self._client._request(
+ "GET", "/nodes", params={"include_all_versions": include_all_versions}
+ )
+
+ def get(self, node_id: str) -> Dict[str, Any]:
+ """Return the current public manifest for a node ID."""
+ return self._client._request("GET", "/nodes/" + segment(node_id))
+
+
+class Graphs(Resource):
+ """Author drafts and discover or publish reusable graph versions."""
+
+ def list(self) -> List[Dict[str, Any]]:
+ """Return the account graph listing."""
+ return self._client._request("GET", "/graphs")
+
+ def get(self, graph_id: str) -> Dict[str, Any]:
+ """Return an accessible graph by ID."""
+ return self._client._request("GET", "/graphs/" + segment(graph_id))
+
+ def validate(self, definition: Dict[str, Any], *, name: Optional[str] = None) -> Dict[str, Any]:
+ """Validate a candidate without saving it or executing nodes.
+
+ An omitted name defaults to the definition name or Untitled Graph.
+ Validation does not guarantee that runtime execution will succeed.
+ """
+ return self._client._request(
+ "POST",
+ "/graphs/validate",
+ json={
+ "name": name if name is not None else definition.get("name", "Untitled Graph"),
+ "definition": definition,
+ },
+ )
+
+ def create(
+ self,
+ *,
+ name: str,
+ definition: Dict[str, Any],
+ description: str = "",
+ tags: Optional[List[str]] = None,
+ ) -> Dict[str, Any]:
+ """Create a public draft and return its server representation.
+
+ Requires graphs:write. Inspect remote state before retrying a request
+ with an uncertain outcome to avoid creating duplicate graphs.
+ """
+ return self._client._request(
+ "POST",
+ "/graphs",
+ json={
+ "name": name,
+ "description": description,
+ "definition": definition,
+ "tags": tags or [],
+ },
+ )
+
+ def update(
+ self, graph_id: str, *, name: str, definition: Dict[str, Any], description: str = ""
+ ) -> Dict[str, Any]:
+ """Replace the submitted draft fields of an owned graph.
+
+ Requires graphs:write. Pass description to preserve its value; the
+ default empty string clears it.
+ """
+ return self._client._request(
+ "PUT",
+ "/graphs/" + segment(graph_id),
+ json={"name": name, "description": description, "definition": definition},
+ )
+
+ def published(self) -> List[Dict[str, Any]]:
+ """Return the published reusable graph catalog."""
+ return self._client._request("GET", "/graphs/published")
+
+ def versions(self, graph_id: str) -> List[Dict[str, Any]]:
+ """Return stored versions of an owned graph."""
+ return self._client._request("GET", "/graphs/" + segment(graph_id) + "/versions")
+
+ def published_version(self, graph_id: str, version: str) -> Dict[str, Any]:
+ """Return a published definition, interface, and composite node manifest."""
+ return self._client._request(
+ "GET", "/graphs/published/" + segment(graph_id) + "/versions/" + segment(version)
+ )
+
+ def interface(self, graph_id: str) -> Dict[str, Any]:
+ """Preview the server-inferred interface of an owned graph draft."""
+ return self._client._request("GET", "/graphs/" + segment(graph_id) + "/interface-preview")
+
+ def publish(
+ self,
+ graph_id: str,
+ *,
+ version: str,
+ interface: Optional[Dict[str, Any]] = None,
+ comment: str = "",
+ ) -> Dict[str, Any]:
+ """Publish an immutable version of an owned graph.
+
+ Requires graphs:publish. The server infers the interface if omitted.
+ This creates public content and is not automatically retried.
+ """
+ payload: Dict[str, Any] = {"version": version, "comment": comment}
+ if interface is not None:
+ payload["interface"] = interface
+ return self._client._request(
+ "POST", "/graphs/" + segment(graph_id) + "/versions", json=payload
+ )
+
+ def execute(
+ self, graph_id: str, *, version: str = "draft", inputs: Optional[Dict[str, Any]] = None
+ ) -> Dict[str, Any]:
+ """Submit an owned graph and return its initial execution state.
+
+ Requires graphs:run and consumes account quota. Select draft, published,
+ or an explicit version, and pass graph inputs as a dictionary. This
+ method does not wait for completion or retry failed submissions.
+ """
+ return self._client._request(
+ "POST",
+ "/graphs/" + segment(graph_id) + "/execute",
+ params={"version": version},
+ json={"inputs": inputs or {}},
+ )
+
+
+class Executions(Resource):
+ """Inspect, cancel, and poll account-owned executions."""
+
+ def get(self, execution_id: str) -> Dict[str, Any]:
+ """Return an owned execution snapshot, including status and outputs."""
+ return self._client._request("GET", "/executions/" + segment(execution_id))
+
+ def cancel(self, execution_id: str) -> Dict[str, Any]:
+ """Request cancellation and return the execution snapshot."""
+ return self._client._request("POST", "/executions/" + segment(execution_id) + "/cancel")
+
+ def wait(
+ self,
+ execution_id: str,
+ *,
+ timeout: float = 120,
+ poll_interval: float = 2,
+ on_update: Optional[Callable[[Dict[str, Any]], None]] = None,
+ ) -> Dict[str, Any]:
+ """Poll until a terminal status or the polling deadline.
+
+ Args:
+ execution_id: Identifier of the execution to observe.
+ timeout: Positive deadline in seconds. Remaining time caps HTTP
+ timeouts, which apply per I/O phase, not to the entire request.
+ poll_interval: Positive delay in seconds between status requests.
+ on_update: Optional synchronous callback for each retrieved snapshot.
+ Callback exceptions propagate to the caller without canceling the run.
+
+ Returns:
+ The terminal snapshot, including failed or canceled executions.
+
+ Raises:
+ ValueError: A timing argument is nonpositive or nonfinite.
+ ExecutionTimeout: Polling expired; the remote run is not canceled.
+ SubforkError: A status request failed; it is not retried.
+ """
+ if (
+ not math.isfinite(timeout)
+ or timeout <= 0
+ or not math.isfinite(poll_interval)
+ or poll_interval <= 0
+ ):
+ raise ValueError("timeout and poll_interval must be positive finite numbers.")
+ deadline = time.monotonic() + timeout
+ while True:
+ remaining = deadline - time.monotonic()
+ if remaining <= 0:
+ raise ExecutionTimeout(execution_id)
+ result = self._client._request(
+ "GET",
+ "/executions/" + segment(execution_id),
+ timeout=min(self._client.timeout, remaining),
+ )
+ if on_update is not None:
+ on_update(result)
+ if result["status"] in {
+ "completed",
+ "failed",
+ "canceled",
+ "cancelled",
+ "outcome_unknown",
+ }:
+ return result
+ time.sleep(min(poll_interval, max(0, deadline - time.monotonic())))
diff --git a/tests/test_cli.py b/tests/test_cli.py
new file mode 100644
index 0000000..58dd78e
--- /dev/null
+++ b/tests/test_cli.py
@@ -0,0 +1,420 @@
+"""Verify CLI contracts using the real SDK with an isolated HTTP transport."""
+
+import io
+import json
+import re
+from pathlib import Path
+from typing import Any
+
+import httpx
+import pytest
+
+from subfork import Subfork, cli
+
+
+@pytest.fixture
+def requests(monkeypatch: pytest.MonkeyPatch) -> list:
+ """Capture requests and provide representative API responses."""
+ seen = []
+
+ def handle(request: httpx.Request) -> httpx.Response:
+ """Return a graph or execution response for each request."""
+ seen.append(request)
+ if request.url.path.endswith("/graphs/g_test"):
+ return httpx.Response(
+ 200, json={"definition": {"name": "Greeting", "nodes": [], "edges": []}}
+ )
+ if request.url.path.endswith("/graphs/validate"):
+ return httpx.Response(200, json={"valid": False, "errors": ["No output"]})
+ return httpx.Response(200, json={"id": "test", "status": "completed"})
+
+ def factory(**kwargs: Any) -> Subfork:
+ """Construct a real client without network access."""
+ return Subfork("synthetic-key", transport=httpx.MockTransport(handle), **kwargs)
+
+ monkeypatch.setattr(cli, "Subfork", factory)
+ return seen
+
+
+def test_export_create_roundtrip(requests: list, tmp_path: Path, capsys: Any) -> None:
+ """Export definitions that create accepts, without leaking response metadata."""
+ output = tmp_path / "graph.json"
+ assert cli.main(["export", "g_test", "--output", str(output)]) == 0
+ definition = json.loads(output.read_text())
+ assert definition == {"name": "Greeting", "nodes": [], "edges": []}
+ assert capsys.readouterr().out == ""
+ assert cli.main(["create", str(output), "--name", "Copy", "--tag", "demo"]) == 0
+ assert json.loads(requests[-1].content) == {
+ "name": "Copy",
+ "definition": definition,
+ "description": "",
+ "tags": ["demo"],
+ }
+ assert json.loads(capsys.readouterr().out)["id"] == "test"
+ assert cli.main(["export", "g_test", "--output", str(output)]) == 1
+ assert json.loads(output.read_text()) == definition
+
+
+@pytest.mark.parametrize("content", ["[]", "invalid", '{"name":'])
+def test_invalid_input_never_sends(
+ requests: list, monkeypatch: pytest.MonkeyPatch, content: str, capsys: Any
+) -> None:
+ """Reject malformed or non-object definitions before any HTTP request."""
+ monkeypatch.setattr("sys.stdin", io.StringIO(content))
+ assert cli.main(["create", "-"]) == 1
+ assert requests == []
+ assert capsys.readouterr().err.startswith("subfork:")
+
+
+def test_publish_and_execute(requests: list, tmp_path: Path) -> None:
+ """Dispatch publication and execution input submission correctly."""
+ assert cli.main(["publish", "g_test", "--version", "v1"]) == 0
+ assert json.loads(requests[-1].content) == {"version": "v1", "comment": ""}
+ inputs = tmp_path / "inputs.json"
+ inputs.write_text('{"greeting": "hello"}')
+ assert cli.main(["execute", "g_test", "--version", "v1", "--inputs", str(inputs)]) == 0
+ assert requests[-1].url.params["version"] == "v1"
+ assert json.loads(requests[-1].content) == {"inputs": {"greeting": "hello"}}
+
+
+def test_validation_failure(requests: list, monkeypatch: pytest.MonkeyPatch, capsys: Any) -> None:
+ """Return a failing exit status while retaining validation diagnostics as JSON."""
+ monkeypatch.setattr("sys.stdin", io.StringIO('{"nodes": []}'))
+ assert cli.main(["validate", "-"]) == 1
+ assert json.loads(capsys.readouterr().out)["valid"] is False
+
+
+def test_help_without_credentials(monkeypatch: pytest.MonkeyPatch) -> None:
+ """Allow help without constructing an authenticated client."""
+ monkeypatch.delenv("SUBFORK_API_KEY", raising=False)
+ with pytest.raises(SystemExit) as result:
+ cli.main(["--help"])
+ assert result.value.code == 0
+
+
+@pytest.mark.parametrize("status", [401, 403])
+def test_api_error_is_safe(monkeypatch: pytest.MonkeyPatch, capsys: Any, status: int) -> None:
+ """Report failed API calls without echoing server bodies or credentials."""
+
+ def factory(**kwargs: Any) -> Subfork:
+ """Create a client whose server returns a sensitive error body."""
+ return Subfork(
+ "synthetic-key",
+ transport=httpx.MockTransport(
+ lambda request: httpx.Response(status, text="sensitive body")
+ ),
+ **kwargs,
+ )
+
+ monkeypatch.setattr(cli, "Subfork", factory)
+ assert cli.main(["list"]) == 1
+ captured = capsys.readouterr()
+ assert captured.out == ""
+ assert str(status) in captured.err
+ if status == 401:
+ assert "SUBFORK_API_KEY" in captured.err
+ assert "does not load .env" in captured.err
+ assert "sensitive" not in captured.err
+ assert "synthetic-key" not in captured.err
+
+
+@pytest.mark.parametrize("command", ["graphs", "nodes", "executions"])
+def test_removed_groups(command: str) -> None:
+ """Reject removed resource groups as usage errors."""
+ with pytest.raises(SystemExit) as result:
+ cli.main([command])
+ assert result.value.code == 2
+
+
+@pytest.mark.parametrize("mode", ["default", "raw", "no-wait", "failed", "timeout"])
+def test_execution_output_modes(monkeypatch: pytest.MonkeyPatch, capsys: Any, mode: str) -> None:
+ """Wait for results by default and preserve explicit diagnostic output modes."""
+ seen = []
+ completed = {
+ "id": "e_test",
+ "status": "completed",
+ "outputs": {"response": ["https://example.com"]},
+ "definition": {"nodes": []},
+ }
+ if mode == "failed":
+ completed["status"] = "failed"
+
+ def handle(request: httpx.Request) -> httpx.Response:
+ """Simulate asynchronous submission followed by completion."""
+ seen.append(request)
+ if request.method == "POST":
+ return httpx.Response(201, json={"id": "e_test", "status": "running", "outputs": {}})
+ return httpx.Response(200, json=completed)
+
+ def factory(**kwargs: Any) -> Subfork:
+ """Construct a client with mocked transport and optional wait timeout."""
+ client = Subfork("synthetic-key", transport=httpx.MockTransport(handle), **kwargs)
+ if mode == "timeout":
+
+ def timeout(*args: Any, **kwargs: Any) -> dict:
+ """Simulate reaching the waiting deadline without canceling."""
+ from subfork import ExecutionTimeout
+
+ raise ExecutionTimeout("e_test")
+
+ monkeypatch.setattr(client.executions, "wait", timeout)
+ return client
+
+ monkeypatch.setattr(cli, "Subfork", factory)
+ arguments = ["execute", "g_test"]
+ if mode in {"raw", "no-wait"}:
+ arguments.append("--" + mode)
+ assert cli.main(arguments) == (1 if mode in {"failed", "timeout"} else 0)
+ captured = capsys.readouterr()
+ assert sum(request.method == "POST" for request in seen) == 1
+ if mode == "timeout":
+ assert captured.out == ""
+ assert "e_test" in captured.err and "not canceled" in captured.err
+ return
+ output = json.loads(captured.out)
+ if mode == "raw":
+ assert output == completed
+ elif mode == "default":
+ assert output == completed["outputs"]
+ else:
+ assert output == {
+ "execution_id": "e_test",
+ "status": "failed" if mode == "failed" else "running",
+ }
+ assert len(seen) == (1 if mode == "no-wait" else 2)
+
+
+@pytest.mark.parametrize("interrupted", [False, True])
+@pytest.mark.parametrize("no_color", [False, True])
+def test_status_cleanup(monkeypatch: pytest.MonkeyPatch, interrupted: bool, no_color: bool) -> None:
+ """Show the graph name and restore terminal output on success or interruption."""
+
+ class Terminal(io.StringIO):
+ """Capture writes as an interactive terminal."""
+
+ def isatty(self) -> bool:
+ """Identify this stream as a terminal."""
+ return True
+
+ terminal = Terminal()
+ monkeypatch.setattr("sys.stderr", terminal)
+ monkeypatch.setenv("TERM", "xterm")
+ monkeypatch.setenv("NO_COLOR", "1" if no_color else "")
+ try:
+ with cli.execution_status("URL Extract Pipe"):
+ expected = "Running" if no_color else "\033[32mRunning\033[0m"
+ assert "Graph URL Extract Pipe .......... " + expected in terminal.getvalue()
+ if interrupted:
+ raise KeyboardInterrupt
+ except KeyboardInterrupt:
+ assert interrupted
+ assert terminal.getvalue().endswith("\r\033[2K")
+
+
+def test_status_redirected_stderr(monkeypatch: pytest.MonkeyPatch) -> None:
+ """Keep redirected progress plain and sanitize control characters in names."""
+ stream = io.StringIO()
+ monkeypatch.setattr("sys.stderr", stream)
+ with cli.execution_status("Greeting\nGraph"):
+ pass
+ assert stream.getvalue() == "Graph Greeting Graph .......... Running\n"
+
+
+def test_active_node_progress(monkeypatch: pytest.MonkeyPatch) -> None:
+ """Animate named parallel nodes and clean up the rendering thread."""
+ import threading
+
+ refreshed = threading.Event()
+
+ class Terminal(io.StringIO):
+ """Capture animated frames and signal when node titles appear."""
+
+ def isatty(self) -> bool:
+ """Enable interactive progress rendering."""
+ return True
+
+ def write(self, value: str) -> int:
+ """Signal a node frame without timing-dependent sleeps."""
+ count = super().write(value)
+ if "Nodes Fetch, Parse" in value:
+ refreshed.set()
+ return count
+
+ terminal = Terminal()
+ monkeypatch.setattr("sys.stderr", terminal)
+ monkeypatch.setenv("TERM", "xterm")
+ monkeypatch.delenv("NO_COLOR", raising=False)
+ with cli.execution_status("Pipeline") as update:
+ update(
+ {
+ "status": "running",
+ "definition": {
+ "nodes": [
+ {"node_instance_id": "a", "title": "Fetch"},
+ {"node_instance_id": "b", "title": "Parse"},
+ ]
+ },
+ "node_executions": {"a": {"status": "running"}, "b": {"status": "running"}},
+ }
+ )
+ assert refreshed.wait(2)
+ assert "\033[33m" in terminal.getvalue()
+ assert "\033[32mRunning\033[0m" in terminal.getvalue()
+ assert not any(thread.name == "subfork-spinner" for thread in threading.enumerate())
+
+
+@pytest.mark.parametrize("terminal", [False, True])
+def test_retained_node_outcomes(
+ monkeypatch: pytest.MonkeyPatch, capsys: Any, terminal: bool
+) -> None:
+ """Retain each terminal node once, including nodes that finish between polls."""
+
+ class ProgressStream(io.StringIO):
+ """Capture either terminal or redirected progress."""
+
+ def isatty(self) -> bool:
+ """Select the requested output mode."""
+ return terminal
+
+ stream = ProgressStream()
+ monkeypatch.setattr("sys.stderr", stream)
+ monkeypatch.setenv("TERM", "xterm")
+ monkeypatch.setenv("NO_COLOR", "1")
+ snapshot = {
+ "status": "failed",
+ "definition": {
+ "nodes": [
+ {"node_instance_id": "a", "title": "Fetch"},
+ {"node_instance_id": "b", "title": "Parse"},
+ {"node_instance_id": "c", "title": "Save"},
+ ]
+ },
+ "node_executions": {
+ "a": {"status": "completed"},
+ "b": {"status": "failed"},
+ "c": {"status": "canceled"},
+ },
+ }
+ with cli.execution_status("Pipeline") as update:
+ update(snapshot)
+ update(snapshot)
+ output = stream.getvalue()
+ for name, outcome in (("Fetch", "Completed"), ("Parse", "Failed"), ("Save", "Canceled")):
+ assert len(re.findall(r"Node " + name + r" \.+ " + outcome + r"\n", output)) == 1
+ assert capsys.readouterr().out == ""
+
+
+def test_node_status_columns_align(monkeypatch: pytest.MonkeyPatch) -> None:
+ """Pad different node names so terminal outcomes start in the same column."""
+ stream = io.StringIO()
+ monkeypatch.setattr("sys.stderr", stream)
+ with cli.execution_status("Demo") as update:
+ update(
+ {
+ "node_executions": {
+ "A": {"status": "completed"},
+ "Longer node name": {"status": "failed"},
+ "Mid": {"status": "canceled"},
+ }
+ }
+ )
+ lines = stream.getvalue().splitlines()[1:]
+ columns = [line.index(state) for line, state in zip(lines, ["Completed", "Failed", "Canceled"])]
+ assert len(columns) == 3
+ assert len(set(columns)) == 1
+
+
+@pytest.mark.parametrize("flag", ["-o", "--out"])
+@pytest.mark.parametrize("mode", [None, "--raw", "--no-wait"])
+def test_execute_result_file(
+ requests: list, tmp_path: Path, capsys: Any, flag: str, mode: Any
+) -> None:
+ """Route results exclusively to a file and reject overwrites before submission."""
+ output = tmp_path / "result.json"
+ arguments = ["execute", "g_test", flag, str(output)]
+ if mode:
+ arguments.append(mode)
+ assert cli.main(arguments) == 0
+ value = json.loads(output.read_text())
+ assert value == (
+ {"id": "test", "status": "completed"}
+ if mode == "--raw"
+ else {"execution_id": "test", "status": "completed"} if mode == "--no-wait" else {}
+ )
+ assert capsys.readouterr().out == ""
+ count = len(requests)
+ assert cli.main(arguments) == 1
+ assert len(requests) == count
+ assert json.loads(output.read_text()) == value
+ captured = capsys.readouterr()
+ assert captured.out == ""
+ assert "already exists" in captured.err
+
+
+@pytest.mark.parametrize(
+ "terminal,no_color,colored", [(True, False, True), (True, True, False), (False, False, False)]
+)
+def test_api_error_display(
+ monkeypatch: pytest.MonkeyPatch, terminal: bool, no_color: bool, colored: bool
+) -> None:
+ """Color only the error code and omit redundant API and program prefixes."""
+
+ class ErrorStream(io.StringIO):
+ """Capture diagnostics with configurable terminal detection."""
+
+ def isatty(self) -> bool:
+ """Return the selected terminal mode."""
+ return terminal
+
+ def factory(**kwargs: Any) -> Subfork:
+ """Raise the SDK's safe conflict explanation."""
+ from subfork import APIError
+
+ raise APIError(
+ "Subfork API returned HTTP 409. A stored graph secret could not be decrypted.",
+ status_code=409,
+ )
+
+ stream = ErrorStream()
+ monkeypatch.setattr("sys.stderr", stream)
+ monkeypatch.setattr(cli, "Subfork", factory)
+ monkeypatch.setenv("TERM", "xterm")
+ monkeypatch.setenv("NO_COLOR", "1" if no_color else "")
+ assert cli.main(["list"]) == 1
+ code = "\033[33m409\033[0m" if colored else "409"
+ assert stream.getvalue() == code + ": A stored graph secret could not be decrypted.\n"
+
+
+@pytest.mark.parametrize("flag", ["-f", "--force"])
+@pytest.mark.parametrize("command,option", [("execute", "-o"), ("export", "--output")])
+def test_force_output(
+ requests: list, tmp_path: Path, capsys: Any, flag: str, command: str, option: str
+) -> None:
+ """Replace the whole output file only when explicitly requested."""
+ output = tmp_path / "existing.json"
+ output.write_text("previous result with extra trailing content")
+ assert cli.main([command, "g_test", option, str(output), flag]) == 0
+ assert json.loads(output.read_text()) == (
+ {} if command == "execute" else {"name": "Greeting", "nodes": [], "edges": []}
+ )
+ assert capsys.readouterr().out == ""
+
+
+def test_force_preserves_output_on_api_failure(
+ monkeypatch: pytest.MonkeyPatch, tmp_path: Path
+) -> None:
+ """Keep the previous result if execution submission fails."""
+
+ def factory(**kwargs: Any) -> Subfork:
+ """Construct a client that receives a server conflict."""
+ return Subfork(
+ "synthetic-key",
+ transport=httpx.MockTransport(lambda request: httpx.Response(409)),
+ **kwargs,
+ )
+
+ monkeypatch.setattr(cli, "Subfork", factory)
+ output = tmp_path / "existing.json"
+ output.write_text("previous result")
+ assert cli.main(["execute", "g_test", "-o", str(output), "--force"]) == 1
+ assert output.read_text() == "previous result"
diff --git a/tests/test_client.py b/tests/test_client.py
new file mode 100644
index 0000000..49b0737
--- /dev/null
+++ b/tests/test_client.py
@@ -0,0 +1,231 @@
+"""Exercise client contracts without network access or real credentials."""
+
+import json
+from typing import Any, Type
+
+import httpx
+import pytest
+
+from subfork import (
+ APIError,
+ AuthenticationError,
+ ExecutionTimeout,
+ PermissionDeniedError,
+ RateLimitError,
+ Subfork,
+ TransportError,
+ ValidationError,
+)
+
+KEY = "synthetic-test-key"
+
+
+def test_bearer_auth_base_url_and_environment(monkeypatch: pytest.MonkeyPatch) -> None:
+ """Verify bearer auth base url and environment."""
+ monkeypatch.setenv("SUBFORK_API_KEY", KEY)
+ seen = []
+
+ def handle(request: httpx.Request) -> httpx.Response:
+ """Handle a mocked HTTP request for this scenario."""
+ seen.append(request)
+ return httpx.Response(200, json=[])
+
+ with Subfork(
+ base_url="http://subfork.localhost", transport=httpx.MockTransport(handle)
+ ) as client:
+ assert client.nodes.list(include_all_versions=True) == []
+ assert str(seen[0].url) == "http://subfork.localhost/api/v1/nodes?include_all_versions=true"
+ assert seen[0].headers["authorization"] == "Bearer " + KEY
+ assert KEY not in str(seen[0].url)
+
+
+@pytest.mark.parametrize(
+ "url",
+ [
+ "http://example.com",
+ "https://user:pass@example.com",
+ "https://example.com/api/v1",
+ "https://example.com?key=x",
+ "http://localhost.evil.example",
+ "https://example.com#fragment",
+ ],
+)
+def test_reject_unsafe_or_ambiguous_origins(url: str) -> None:
+ """Verify reject unsafe or ambiguous origins."""
+ with pytest.raises(ValueError):
+ Subfork(KEY, base_url=url)
+
+
+def test_credentials_required(monkeypatch: pytest.MonkeyPatch) -> None:
+ """Verify credentials required."""
+ monkeypatch.delenv("SUBFORK_API_KEY", raising=False)
+ with pytest.raises(ValueError):
+ Subfork()
+
+
+@pytest.mark.parametrize(
+ "status,error",
+ [
+ (401, AuthenticationError),
+ (403, PermissionDeniedError),
+ (409, APIError),
+ (422, ValidationError),
+ (429, RateLimitError),
+ (503, APIError),
+ ],
+)
+def test_errors_do_not_echo_response_secrets(status: int, error: Type[APIError]) -> None:
+ """Verify errors do not echo response secrets."""
+ with Subfork(
+ KEY,
+ transport=httpx.MockTransport(
+ lambda request: httpx.Response(
+ status, json={"detail": KEY}, headers={"Retry-After": "12"}
+ )
+ ),
+ ) as client:
+ with pytest.raises(error) as caught:
+ client.graphs.list()
+ assert caught.value.status_code == status
+ assert caught.value.retry_after == "12"
+ assert KEY not in str(caught.value)
+
+
+def test_redirects_and_failed_writes_are_not_retried() -> None:
+ """Verify redirects and failed writes are not retried."""
+ requests = []
+
+ def handle(request: httpx.Request) -> httpx.Response:
+ """Handle a mocked HTTP request for this scenario."""
+ requests.append(request)
+ return httpx.Response(307, headers={"Location": "https://other.example"})
+
+ with Subfork(KEY, transport=httpx.MockTransport(handle)) as client:
+ with pytest.raises(APIError):
+ client.graphs.create(name="Test", definition={"name": "Test"})
+ assert len(requests) == 1
+
+
+def test_authoring_requests_match_server_contract() -> None:
+ """Verify authoring requests match server contract."""
+ seen = []
+
+ def handle(request: httpx.Request) -> httpx.Response:
+ """Handle a mocked HTTP request for this scenario."""
+ seen.append(
+ (
+ request.method,
+ request.url.path,
+ dict(request.url.params),
+ json.loads(request.content) if request.content else None,
+ )
+ )
+ return httpx.Response(200, json={"id": "g_test"})
+
+ with Subfork(KEY, transport=httpx.MockTransport(handle)) as client:
+ definition = {"name": "Demo", "nodes": [], "edges": []}
+ client.graphs.validate(definition)
+ client.graphs.create(name="Demo", definition=definition)
+ client.graphs.update("g_test", name="Demo", definition=definition)
+ client.graphs.publish("g_test", version="v1", interface={"inputs": {}, "outputs": {}})
+ client.graphs.execute("g_test", version="v1", inputs={"value": 3})
+ assert [row[:2] for row in seen] == [
+ ("POST", "/api/v1/graphs/validate"),
+ ("POST", "/api/v1/graphs"),
+ ("PUT", "/api/v1/graphs/g_test"),
+ ("POST", "/api/v1/graphs/g_test/versions"),
+ ("POST", "/api/v1/graphs/g_test/execute"),
+ ]
+ assert seen[0][3] == {"name": "Demo", "definition": definition}
+ assert seen[-2][3]["interface"] == {"inputs": {}, "outputs": {}}
+ assert seen[-1][2:] == ({"version": "v1"}, {"inputs": {"value": 3}})
+
+
+def test_transport_failure_is_not_retried() -> None:
+ """Verify transport failure is not retried."""
+ seen = []
+
+ def handle(request: httpx.Request) -> httpx.Response:
+ """Handle a mocked HTTP request for this scenario."""
+ seen.append(request)
+ raise httpx.ReadTimeout(KEY)
+
+ with Subfork(KEY, transport=httpx.MockTransport(handle)) as client:
+ with pytest.raises(TransportError) as caught:
+ client.graphs.execute("g_test")
+ assert len(seen) == 1 and KEY not in str(caught.value)
+
+
+@pytest.mark.parametrize("status", ["completed", "failed", "canceled", "outcome_unknown"])
+def test_wait_returns_terminal_result(status: str) -> None:
+ """Verify wait returns terminal result."""
+ with Subfork(
+ KEY,
+ transport=httpx.MockTransport(
+ lambda request: httpx.Response(200, json={"id": "e_test", "status": status})
+ ),
+ ) as client:
+ assert client.executions.wait("e_test")["status"] == status
+
+
+def test_wait_timeout_does_not_cancel(monkeypatch: pytest.MonkeyPatch) -> None:
+ """Verify wait timeout does not cancel."""
+ from subfork import resources
+
+ ticks = iter([0, 0, 1, 2])
+ monkeypatch.setattr(resources.time, "monotonic", lambda: next(ticks))
+ monkeypatch.setattr(resources.time, "sleep", lambda seconds: None)
+ seen = []
+
+ def handle(request: httpx.Request) -> httpx.Response:
+ """Handle a mocked HTTP request for this scenario."""
+ seen.append(request)
+ return httpx.Response(200, json={"status": "running"})
+
+ with Subfork(KEY, transport=httpx.MockTransport(handle)) as client:
+ with pytest.raises(ExecutionTimeout):
+ client.executions.wait("e_test", timeout=1)
+ assert len(seen) == 1 and seen[0].method == "GET"
+
+
+def test_wait_reports_snapshots() -> None:
+ """Deliver both active and terminal snapshots to an optional progress callback."""
+ snapshots = []
+ statuses = iter(["running", "completed"])
+ with Subfork(
+ "synthetic-key",
+ transport=httpx.MockTransport(
+ lambda request: httpx.Response(200, json={"status": next(statuses)})
+ ),
+ ) as client:
+ result = client.executions.wait("e_test", poll_interval=0.001, on_update=snapshots.append)
+ assert [snapshot["status"] for snapshot in snapshots] == ["running", "completed"]
+ assert result == snapshots[-1]
+
+
+@pytest.mark.parametrize(
+ "detail,recognized",
+ [
+ (
+ "Graph secret '"
+ + KEY
+ + "' could not be decrypted. The server encryption key may have changed; re-enter this secret in Graph Settings.",
+ True,
+ ),
+ (KEY, False),
+ ({"secret": KEY}, False),
+ (None, False),
+ ],
+)
+def test_decryption_error_guidance(detail: Any, recognized: bool) -> None:
+ """Explain known conflicts without echoing even the server-provided secret name."""
+ with Subfork(
+ KEY,
+ transport=httpx.MockTransport(lambda request: httpx.Response(409, json={"detail": detail})),
+ ) as client:
+ with pytest.raises(APIError) as caught:
+ client.graphs.execute("g_test")
+ message = str(caught.value)
+ assert caught.value.status_code == 409
+ assert KEY not in message
+ assert ("Graph Settings > Secrets" in message) is recognized
diff --git a/tests/test_version.py b/tests/test_version.py
new file mode 100644
index 0000000..c31decf
--- /dev/null
+++ b/tests/test_version.py
@@ -0,0 +1,56 @@
+"""Exercise release targets in isolated copies without bumping this checkout."""
+
+import shutil
+import subprocess
+from pathlib import Path
+
+import pytest
+
+ROOT = Path(__file__).resolve().parents[1]
+
+
+@pytest.mark.parametrize(
+ "target,version",
+ [
+ ("bump-patch", "2.3.5"),
+ ("bump-minor", "2.4.0"),
+ ("bump-major", "3.0.0"),
+ ("bump-version", "4.5.6"),
+ ("version", "2.3.4"),
+ ("check-version", "2.3.4"),
+ ],
+)
+def test_version_targets(tmp_path: Path, target: str, version: str) -> None:
+ """Run Make targets against a fixture and preserve unrelated metadata."""
+ import sys
+
+ (tmp_path / "scripts").mkdir()
+ (tmp_path / "src/subfork").mkdir(parents=True)
+ shutil.copy(ROOT / "Makefile", tmp_path)
+ shutil.copy(ROOT / "scripts/bump-version.py", tmp_path / "scripts")
+ project = tmp_path / "pyproject.toml"
+ init = tmp_path / "src/subfork/__init__.py"
+ project.write_text('[project]\nname = "subfork"\nversion = "2.3.4"\n')
+ init.write_text('__version__ = "2.3.4"\n')
+ result = subprocess.run(
+ ["make", target, "PYTHON=" + sys.executable, "VERSION=4.5.6"],
+ cwd=tmp_path,
+ capture_output=True,
+ text=True,
+ )
+ assert result.returncode == 0, result.stderr
+ assert version in result.stdout
+ assert project.read_text() == '[project]\nname = "subfork"\nversion = "' + version + '"\n'
+ assert init.read_text() == '__version__ = "' + version + '"\n'
+ # A mismatch must fail before either file is changed.
+ init.write_text('__version__ = "9.9.9"\n')
+ before = project.read_text()
+ result = subprocess.run(
+ [sys.executable, "scripts/bump-version.py", "--bump", "patch"],
+ cwd=tmp_path,
+ capture_output=True,
+ text=True,
+ )
+ assert result.returncode == 1
+ assert project.read_text() == before
+ assert init.read_text() == '__version__ = "9.9.9"\n'