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 + +
+ +# 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](assets/subfork-banner.png) + +# 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'