Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion .github/workflows/images.yml
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ on:
- "hello-world/**"
- "echo/**"
- "tiles/**"
- "api-proxy/**"
- ".github/workflows/images.yml"
pull_request:
paths: *image_paths
Expand All @@ -34,7 +35,7 @@ jobs:
strategy:
fail-fast: false
matrix:
example: [hello-world, echo, tiles]
example: [hello-world, echo, tiles, api-proxy]
steps:
- uses: actions/checkout@v7

Expand Down
21 changes: 11 additions & 10 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ flowchart LR

The orchestrator is a **transparent reverse proxy**: every endpoint you expose is passed through to your app unchanged, so you write an ordinary service and it runs on the network as-is. The transports supported today:

- **HTTP** request/response — the common case. (`hello-world`, `tiles`)
- **HTTP** request/response — the common case. (`hello-world`, `tiles`, `api-proxy`)
- **HTTP + SSE** — streamed / token responses. (`vllm`)
- **Trickle** — continuous realtime video in/out. (`echo`)
- **WebSocket** — long-lived bidirectional sessions. (external: `scope`)
Expand All @@ -37,12 +37,13 @@ Need a schema that isn't here? [Open an issue](https://github.com/livepeer/runne

## Examples

| Example | Goal | Registration | Mode | Transport | Pricing |
| ------------------------------ | ----------------------------------------------- | ------------ | ---------------------------------- | ----------------- | ----------------- |
| [`hello-world`](./hello-world) | The simplest app: one request, one response | dynamic | single-shot | HTTP (JSON) | fixed |
| [`tiles`](./tiles) | Capacity fan-out — one call per tile | dynamic | single-shot | HTTP (base64 PNG) | fixed |
| [`echo`](./echo) | Realtime video, transformed and echoed back | dynamic | persistent | trickle | — (offchain only) |
| [`vllm`](./vllm) | Drop-in OpenAI API; the client stays unmodified | static | persistent (single-shot by nature) | HTTP + SSE | hour |
| Example | Goal | Registration | Mode | Transport | Pricing |
| ------------------------------ | ---------------------------------------------------------------------- | ------------ | ---------------------------------- | -------------------- | ----------------- |
| [`hello-world`](./hello-world) | The simplest app: one request, one response | dynamic | single-shot | HTTP (JSON) | fixed |
| [`tiles`](./tiles) | Capacity fan-out — one call per tile | dynamic | single-shot | HTTP (base64 PNG) | fixed |
| [`api-proxy`](./api-proxy) | Resell any HTTP API — the operator holds the key, callers pay per call | dynamic | single-shot | HTTP (JSON envelope) | fixed |
| [`echo`](./echo) | Realtime video, transformed and echoed back | dynamic | persistent | trickle | — (offchain only) |
| [`vllm`](./vllm) | Drop-in OpenAI API; the client stays unmodified | static | persistent (single-shot by nature) | HTTP + SSE | hour |

Start with `hello-world` (the smallest end-to-end path); the others each layer on one new idea. More will follow, including a full example that exercises every feature. Each is self-contained and runs **offchain** (free, no wallet); most also run **on-chain** (paid) — see each README.

Expand All @@ -52,7 +53,7 @@ This set stays **minimal and curated**: it covers each value of the axes above (

How the app attaches to the orchestrator:

- **Dynamic** — the app self-registers via the SDK (`register_runner`) and heartbeats; the orchestrator drops it when heartbeats stop. Best for apps that come and go. (`hello-world`, `echo`)
- **Dynamic** — the app self-registers via the SDK (`register_runner`) and heartbeats; the orchestrator drops it when heartbeats stop. Best for apps that come and go. (`hello-world`, `api-proxy`, `echo`)
- **Static** — the orchestrator is configured with the app's URL in a `runners.json` and health-polls it; the app needs no SDK. Best for fixed, long-running deployments. (`vllm`)

The arrow flips — dynamic, the app announces itself; static, the orchestrator is told about a passive app:
Expand All @@ -74,7 +75,7 @@ flowchart LR
Chosen _at_ registration (above); **defaults to `persistent`** — set on both `register_runner(...)` and in `runners.json`. The examples set it explicitly.

- **Persistent** — a held-open session the client reserves and releases, billed per second of wall-clock (or once, with fixed pricing). Best for realtime / streaming. (`echo`, `vllm`)
- **Single-shot** — one request in, one response out; the orchestrator reserves a session per call and releases it when the response returns, so the client manages no session at all. Best for batch / request-response. (`hello-world`, `tiles`)
- **Single-shot** — one request in, one response out; the orchestrator reserves a session per call and releases it when the response returns, so the client manages no session at all. Best for batch / request-response. (`hello-world`, `tiles`, `api-proxy`)

> [!NOTE]
> The `vllm` example is single-shot by nature but stays **persistent** for now: it meters per second across a reserved session, and true per-token billing is brokerage for the gateway/signer layer.
Expand All @@ -83,7 +84,7 @@ Chosen _at_ registration (above); **defaults to `persistent`** — set on both `

The client side depends on the runner's mode:

- **Single-shot** — **discover → call**: find the app via `runner_selector`, then one `call_runner`. The orchestrator reserves a session for the call and releases it when the response returns; on the paid path `call_runner` answers the 402 payment challenge inline. (`hello-world`, `tiles`)
- **Single-shot** — **discover → call**: find the app via `runner_selector`, then one `call_runner`. The orchestrator reserves a session for the call and releases it when the response returns; on the paid path `call_runner` answers the 402 payment challenge inline. (`hello-world`, `tiles`, `api-proxy`)
- **Persistent** — **discover → reserve → call → release**: reserve a session (`reserve_session`), call it — `call_runner`, streamed frames, or a WebSocket, depending on transport — then release it (`stop_runner_session`), which settles payment on-chain. (`echo`, `vllm`)

Each example's `client.py` shows its exact calls — grep `# Livepeer:` to find them.
Expand Down
29 changes: 29 additions & 0 deletions api-proxy/.env.example
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
# Copy to .env (gitignored) and fill in. Never commit secrets.
# Keystore dirs: absolute paths OUTSIDE this repo, mounted read-only.

# Upstream credential the app injects (required, offchain too):
# huggingface.co → settings → tokens.
HF_TOKEN=hf_your_token

NETWORK=arbitrum-one-mainnet
ETH_RPC_URL=https://arb1.arbitrum.io/rpc

# Signer (payer): needs an on-chain deposit + reserve.
SIGNER_KEYSTORE_DIR=/absolute/path/to/signer-keystore
SIGNER_ETH_ACCT=0xYourSignerAddress
SIGNER_ETH_PASSWORD=your-signer-keystore-password

# Orchestrator operating key (split-key): needs ETH for gas to redeem tickets.
ORCH_KEYSTORE_DIR=/absolute/path/to/operator-keystore
ORCH_ETH_ACCT=0xYourOperatorAddress
ORCH_ETH_PASSWORD=your-operator-keystore-password
# Registered orch = ticket recipient (-ethOrchAddr); empty = use the operating key.
ORCH_ONCHAIN_ADDR=0xYourRegisteredOrchestrator

# Runner price (on-chain): USD billed once per call (fixed pricing).
# Keep under ~0.00019: the signer signs at most 100 tickets per payment,
# and the demo orchestrator runs -ticketEV=1e9 (numTickets = fee / ticketEV).
PRICE=0.0001
# Signer's max-price cap (payer side), compared per billing unit. With fixed
# pricing the unit is one call, so this must exceed PRICE.
MAX_PRICE_PER_UNIT=0.000111USD
2 changes: 2 additions & 0 deletions api-proxy/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
# Client output
api-proxy-out.jpg
20 changes: 20 additions & 0 deletions api-proxy/Dockerfile
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
# api-proxy example app (http server).
FROM python:3.12-slim

# Flush stdout/stderr immediately so output isn't block-buffered in `docker logs`.
ENV PYTHONUNBUFFERED=1

RUN apt-get update \
&& apt-get install -y --no-install-recommends git \
&& rm -rf /var/lib/apt/lists/*

# livepeer-gateway SDK isn't on PyPI yet; install from Git.
RUN pip install --no-cache-dir \
"livepeer-gateway @ git+https://github.com/livepeer/livepeer-python-gateway@ja/live-runner"
Comment thread
rickstaa marked this conversation as resolved.
Outdated

WORKDIR /app
COPY runner.py client.py ./

EXPOSE 8989

ENTRYPOINT ["python", "runner.py"]
59 changes: 59 additions & 0 deletions api-proxy/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
# API-proxy app (resell an upstream API)

Wraps an existing HTTP API so it's reachable, and payable, through the Livepeer network. `POST /proxy` takes a JSON envelope describing the upstream call and returns the upstream response — the app runs no model and knows nothing about what it forwards. The demo upstream is the **Hugging Face text-to-image inference API** ([Stable Diffusion 3 medium](https://huggingface.co/stabilityai/stable-diffusion-3-medium-diffusers) by default), but `--upstream` points it at any REST API.

| | |
| ------------ | ---------------------------------------- |
| App id | `livepeer-example/api-proxy` |
| Runner mode | single-shot |
| Registration | dynamic (self-registers via the SDK) |
| Transport | HTTP (JSON envelope in, JSON/base64 out) |
| Port | 8989 |

Prerequisites (Docker, `uv`, and the not-yet-released `livepeer-gateway` SDK — pinned in `pyproject.toml`) and the shared on-chain/payment setup live in the [repo README](../README.md). The demo upstream additionally needs a **Hugging Face API token** (`HF_TOKEN`, from [huggingface.co → settings → tokens](https://huggingface.co/settings/tokens)) with inference-provider credits.

## How it's wired

The app is **dynamically registered**: it self-registers with the orchestrator via `register_runner` ([runner.py](runner.py)) and exposes a single `POST /proxy`, reverse-proxied through the orchestrator. Each call forwards `{"method", "path", "headers", "json"}` to `<upstream>/<path>` and returns `{"status", "headers", "body"}` for text upstream bodies or `{"status", "headers", "body_b64"}` for binary ones (a generated image, say). The client calls it with `runner_selector` → `call_runner` ([client.py](client.py)) — discover, then one **single-shot** call per request. There is no session to manage: the orchestrator reserves one per call and releases it when the response returns. Grep `# Livepeer:` in either file to see the exact calls.

## Proxying an API — what this shows

Most real apps don't host models; they call an API. This example shows that a runner can be exactly that call: the same thin proxy you would deploy anywhere, registered on the network unchanged.

The interesting part is **who holds the key**. The upstream credential (`UPSTREAM_TOKEN`) lives with the runner operator; the app injects it as a Bearer token on every forward and drops any `Authorization` a caller sends. Callers never see an API key — they discover the app and pay **per call through Livepeer**, while the operator pays the upstream and sets `PRICE` above the per-call upstream cost. **Fixed pricing** is the natural fit: one call is one bounded unit of work, so the runner bills one flat price per call instead of metering time (compare [`vllm`](../vllm), where open-ended sessions make per-second metering the better fit).

## Run offchain (free)

```sh
HF_TOKEN=hf_... docker compose up -d --build
curl -sk https://localhost:8935/discovery | jq '.[].runners[].app' # confirm livepeer-example/api-proxy registered
uv run client.py --prompt "a watercolor painting of a llama writing code"
docker compose down
```

`compose.yml` brings up an orchestrator (`-useLiveRunners`) and the app (proxying `https://router.huggingface.co`). The client builds the envelope for one text-to-image call, sends it through the orchestrator, and writes `api-proxy-out.jpg`.

## Run on-chain (paid)

Layer `compose.onchain.yml` to run the orchestrator on-chain with a remote signer paying each call — one fixed payment per image. For the required RPC and wallets see [On-chain (paid) setup](../README.md#on-chain-paid-setup) in the repo README.

```sh
cp .env.example .env # fill in HF_TOKEN, RPC, network, keystore paths, accounts, pricing
docker compose -f compose.yml -f compose.onchain.yml up -d --build
uv run client.py --prompt "a watercolor painting of a llama writing code" \
--discovery https://localhost:8935/discovery \
--signer http://localhost:7936
docker compose -f compose.yml -f compose.onchain.yml down
```

Each call is one paid single-shot session — the orchestrator reserves it, takes one fixed payment, and releases it when the response returns.

## Run without Docker

Start an orchestrator built from go-livepeer `v0.9.0` or newer (see [Build from source](https://docs.livepeer.org/v1/orchestrators/guides/install-go-livepeer#build-from-source)), then the app and client directly:

```sh
./livepeer -orchestrator -useLiveRunners -serviceAddr localhost:8935 -orchSecret abcdef -v 6
UPSTREAM_TOKEN=hf_... uv run runner.py --orchestrator https://localhost:8935 --orchSecret abcdef
uv run client.py --prompt "a watercolor painting of a llama writing code"
```
97 changes: 97 additions & 0 deletions api-proxy/client.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
#!/usr/bin/env python3
"""api-proxy client: discover a runner, proxy a text-to-image call, save the image.

Builds the generic /proxy envelope for one concrete upstream — the Hugging Face
text-to-image inference API — and decodes the binary response. Any other REST
API is the same envelope with a different method/path/json.

Livepeer integration (grep `# Livepeer:`):
1. runner_selector() — discover orchestrators advertising the app
2. call_runner() — call the app through the orchestrator; on the paid path it
answers the 402 payment challenge inline (one fixed payment
per call). Single-shot needs no reserve/stop: the
orchestrator reserves the session for this one request and
releases it when the response returns.
"""

from __future__ import annotations

import argparse
import asyncio
import base64
import logging
from pathlib import Path

from livepeer_gateway.errors import LivepeerGatewayError
from livepeer_gateway.live_runner import call_runner
from livepeer_gateway.selection import runner_selector

DEFAULT_DISCOVERY = "https://localhost:8935/discovery"
APP_ID = "livepeer-example/api-proxy"
DEFAULT_MODEL = "stabilityai/stable-diffusion-3-medium-diffusers"
DEFAULT_OUTPUT = "api-proxy-out.jpg"

log = logging.getLogger("api-proxy-client")


def _parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description="Run the api-proxy Live Runner demo.")
parser.add_argument(
"--prompt", default="a watercolor painting of a llama writing code"
)
parser.add_argument(
"--model",
default=DEFAULT_MODEL,
help="Hugging Face text-to-image model (the upstream path).",
)
parser.add_argument("--output", default=DEFAULT_OUTPUT, help="output image path")
parser.add_argument("--discovery", default=DEFAULT_DISCOVERY)
parser.add_argument(
"--signer", default="", help="Remote signer base URL (on-chain/paid path)."
)
return parser.parse_args()


async def main() -> None:
logging.basicConfig(
level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s"
)
args = _parse_args()
try:
cursor = await runner_selector( # Livepeer: 1
discovery_url=args.discovery, app=APP_ID
)
runner = cursor.candidates[0]
log.info("app_url=%s", runner.url)

# The envelope the app forwards upstream: here a Hugging Face
# text-to-image call, but any method/path/json works.
envelope = {
"method": "POST",
"path": f"/hf-inference/models/{args.model}",
"json": {"inputs": args.prompt},
}
result = await call_runner( # Livepeer: 2
runner=runner, # discovery metadata tells call_runner the price unit
runner_url=runner.url.rstrip("/") + "/proxy",
payload=envelope,
signer_url=args.signer.strip() or None,
timeout=120.0, # a hosted diffusion model can take tens of seconds
)

status = result.data.get("status")
if status != 200:
body = result.data.get("body") or result.data.get("error")
raise LivepeerGatewayError(f"upstream returned {status}: {body}")
b64 = result.data.get("body_b64")
if not isinstance(b64, str) or not b64:
raise LivepeerGatewayError("upstream response was not binary (no body_b64)")
out_path = Path(args.output).expanduser()
out_path.write_bytes(base64.b64decode(b64))
log.info("wrote %s", out_path)
except LivepeerGatewayError as exc:
raise SystemExit(f"ERROR: {exc}") from exc


if __name__ == "__main__":
asyncio.run(main())
34 changes: 34 additions & 0 deletions api-proxy/compose.onchain.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
# On-chain payment overlay for api-proxy. Layer it on the offchain base:
# docker compose -f compose.yml -f compose.onchain.yml up -d --build
#
# Adds the shared remote signer, re-points the orchestrator on-chain (see
# ../compose.onchain.yml), and registers the app with a price so the
# orchestrator issues a payment challenge. Requires a local .env (gitignored);
# copy .env.example and fill it in. Then pay through the signer:
# uv run client.py --prompt "..." \
# --discovery https://localhost:8935/discovery \
# --signer http://localhost:7936

services:
signer:
extends:
file: ../compose.onchain.yml
service: signer
ports:
- "7936:7936"

orchestrator:
extends:
file: ../compose.onchain.yml
service: orchestrator

# Re-declare the command to advertise a price (base file registers free).
app:
command:
- --host=0.0.0.0
- --orchestrator=https://orchestrator:8935
- --orchSecret=abcdef
- --runner-url=http://app:8989
- --upstream=https://router.huggingface.co
# Billed once per call (fixed pricing); price cap in .env.example.
- --price=${PRICE}
31 changes: 31 additions & 0 deletions api-proxy/compose.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
# End-to-end offchain demo: orchestrator + api-proxy app.
#
# The orchestrator service is defined once in ../compose.orchestrator.yml and
# pulled in with `extends`; this file only adds the app. Needs HF_TOKEN (the
# upstream credential the app injects) from the shell or a local .env. Once up,
# call it from the host with the SDK:
# HF_TOKEN=hf_... docker compose up -d --build
# uv run client.py --prompt "..." --discovery https://localhost:8935/discovery

services:
orchestrator:
extends:
file: ../compose.orchestrator.yml
service: orchestrator

app:
build: .
container_name: example_apps_api_proxy
environment:
# The operator-held upstream credential (huggingface.co → settings → tokens).
UPSTREAM_TOKEN: ${HF_TOKEN:?set HF_TOKEN in the shell or .env}
# Wait for the orchestrator's healthcheck so registration doesn't race its boot.
depends_on:
orchestrator:
condition: service_healthy
command:
- --host=0.0.0.0
- --orchestrator=https://orchestrator:8935
- --orchSecret=abcdef
- --runner-url=http://app:8989
- --upstream=https://router.huggingface.co
Loading
Loading