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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 34 additions & 0 deletions .github/workflows/tests.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
name: tests
on:
pull_request:
push:
branches: [main]
jobs:
test:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-python@v5
with:
python-version: '3.11'
- run: pip install -e '.[dev]' matplotlib
- name: Start isolated accounting database
run: |
docker run -d --rm --name ow-cost-tests -e POSTGRES_PASSWORD=ow-test-only postgres:16
for attempt in $(seq 1 30); do
if docker exec ow-cost-tests pg_isready -U postgres; then exit 0; fi
sleep 1
done
exit 1
- name: Unit, API and database tests
env:
OW_COST_TEST_CONTAINER: ow-cost-tests
run: pytest tests --ignore=tests/test_integration.py -q
- uses: actions/setup-node@v4
with:
node-version: '22'
cache: npm
cache-dependency-path: openweights/dashboard/frontend/package-lock.json
- name: Build dashboard
working-directory: openweights/dashboard/frontend
run: npm ci --no-audit --no-fund && npm run build
10 changes: 9 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,14 @@ This repo is research code. Please use github issues or contact me via email (ni
An openai-like sdk with the flexibility of working on a local GPU: finetune, inference, API deployments and custom workloads on managed runpod instances.


## Cost tracking

View estimated compute spending by job, worker, API key and user in the **Costs**
dashboard or `ow.costs.report()`. Startup and idle overhead are shown separately and
included equally across a worker's jobs by default. Organization admins can set
lifetime per-key spending limits. [Accounting rules, limits and required database
migration](docs/cost-tracking.md).

## Installation
Run `pip install openweights` or install from source via `pip install -e .`

Expand Down Expand Up @@ -68,7 +76,7 @@ class MyCustomJob(Jobs):
}
params: Type[BaseModel] = MyParams # Your Pydantic model for params
requires_vram_gb: int = 24
base_image: str = 'nielsrolf/ow-unsloth:v0.12.1' # optional
base_image: str = 'nielsrolf/ow-unsloth:v0.13.0' # optional

def get_entrypoint(self, validated_params: BaseModel) -> str:
# Get the entrypoint command for the job.
Expand Down
134 changes: 134 additions & 0 deletions docs/cost-tracking.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,134 @@
# Cost tracking and spending limits

OpenWeights 0.13 adds organization-scoped compute estimates in USD. Open **Costs**
in the dashboard, or use:

```python
from openweights import OpenWeights
ow = OpenWeights()
report = ow.costs.report()
print(report['total_usd'])
for job in report['jobs']:
print(job['job_id'], job['direct_usd'], job['overhead_usd'], job['total_usd'])
```

Reports include jobs, workers, API keys, submitting users, configured limits, and
an `as_of` timestamp. Dashboard totals refresh every 30 seconds. Reports cover the
entire recorded lifetime; there is no monthly reset or date-window filter. Job and
worker arrays are paginated (100 rows by default); use `ow.costs.report(limit=100,
offset=100)` for the next page. Totals and key/user summaries always cover all rows.

## Accounting rules

- The manager snapshots the pod's `costPerHr` from RunPod at provisioning. If
unavailable, it records the existing hardware estimate multiplied by GPU count.
Unknown hardware remains **unpriced**, not free. The source is visible per worker.
Provider pod rates already cover the whole pod and are not multiplied again.
- Worker cost is elapsed time from the successful provisioning attempt's start to
confirmed pod termination, multiplied by the saved hourly rate. This includes
startup, downloads, execution, uploads and idle time. Shutdown requests do not
stop the clock. A failed termination is retried by the manager.
- Direct job cost sums execution intervals across **all runs**, including failures,
cancellations and retries. Run completion timestamps freeze when status first
leaves `in_progress`; later log uploads do not extend execution time.
- Worker overhead is its total cost minus direct execution cost. It is retained
separately and divided **equally among distinct jobs on that worker**, not among
runs and not in proportion to runtime. A retry does not buy another overhead share.
- Default job/key/user totals include this allocated overhead. Toggle it off in
the dashboard to inspect execution only. A worker that never executes a job has
unallocated overhead, visible separately in the organization total.
- Example: a worker costs $8, job A executes for $3 (over two attempts), and job B
for $1. The $4 overhead adds $2 to each job: A totals $5 and B totals $3.
- Allocations change until worker termination and can decrease when another job
shares the same worker. Organization totals count each worker once.
- Jobs retain the original submitting key and the key's creator as the submitting
user. Signed-in user submissions use their authenticated user ID. Deduplicated
jobs and retries keep the original attribution, even if another key restarts
them. Legacy organization JWTs and old jobs remain unattributed.
- Separate accounting records survive deletion of jobs, runs, API keys and users.
Deleted job IDs with accounting history cannot be reused. Run identity is fixed;
retry by creating another run, not reopening or reassigning an old run.

These are operational estimates, not invoice reconciliation. They do not separately
meter network volumes, external APIs, storage, taxes or credits. Hardware fallback
prices are indicative. Pre-upgrade workers have unknown prices; historical usage is
not retroactively priced. Unknown runs and workers produce an incomplete-total
warning. Non-managed/local runs have no inferred hardware price. Historical closed
runs use the best available `updated_at` timestamp during migration.

## API key spending limits

A signed-in organization **admin user** can choose an API key in the Costs page and
save a lifetime USD limit. `0` blocks work; blank removes the limit. API-key logins
can view costs but cannot edit budgets, even though older organization APIs treat
API keys as administrators.

For automation with an administrator's Supabase user JWT:

```python
admin = OpenWeights(auth_token=admin_user_jwt, organization_id=organization_id)
admin.costs.set_limit(api_token_id, 25)
admin.costs.set_limit(api_token_id, None) # remove limit
```

The REST equivalents are `GET /organizations/{org_id}/costs` and
`PUT /organizations/{org_id}/costs/limits/{token_id}` with
`{"amount_usd": 25}` (or `null` to remove). Amounts must be finite and nonnegative.

Budgets include direct execution and allocated overhead. Database triggers reject
new submissions, restarts, acquisition and run creation after the threshold is
reached. The manager checks every approximately 15 seconds and cancels pending and
running jobs for exhausted keys. Workers follow their existing cancellation path
(5-second polling plus up to about 60 seconds of log-flush delay). Idle termination
can add further overhead. Budgeted jobs cannot start on unpriced workers.

**These are accrued-spend limits, not strict prepaid caps.** Parallel jobs may start
before a threshold is reached, startup cost may only be allocated when a worker
first runs a job, and shutdown delays, manager outages and later overhead can cause
overspend. Limits do not reserve the estimated maximum price of pending jobs.
Previously canceled jobs require an explicit restart after increasing a limit.
Limits attach to individual keys, not to all credentials held by a person; this
feature does not turn the existing organization-admin API keys into untrusted,
restricted credentials. Use independently managed keys and trusted organization
members. Unallocated overhead cannot be assigned to a key or user.

## Rollout

1. Back up the database and apply
`supabase/migrations/20260922000000_cost_tracking.sql` through your existing
Supabase migration process (or SQL editor). It is transactional and must run
exactly once. It adds columns, accounting tables, protected RPCs and triggers;
it does not reprice old usage. Do not apply the original schema to production.
With Supabase CLI management login, an explicit project can be migrated using
`supabase db query --linked --project-ref YOUR_PROJECT_REF --file supabase/migrations/20260922000000_cost_tracking.sql`.
This runs SQL directly; record the migration in your deployment history if your
environment also uses `supabase db push`.
2. Upgrade the manager/dashboard and managed worker images to v0.13.0, and upgrade
SDK clients with `pip install --upgrade openweights==0.13.0`. Run the migration
**before** the new manager starts: it deliberately fails closed if the budget
enforcement RPC is unavailable.
3. Allow existing workers to drain, then terminate them and provision new workers
to start recording rates. Existing workers are marked unpriced. Older workers
still benefit from database triggers; the updated manager confirms termination.
4. Confirm a test job appears under Costs with its submitting key/user, inspect
the worker rate source, and exercise a zero-dollar limit with a test key.

Keep the additive schema when rolling application code back; preserve the accounting
tables. Older managers do not perform periodic cancellation or price snapshots, so
budget enforcement is degraded on rollback. No production database is modified by
package installation or the test suite.

## Validation

```sh
docker run -d --rm --name ow-cost-tests -e POSTGRES_PASSWORD=ow-test-only postgres:16
OW_COST_TEST_CONTAINER=ow-cost-tests .venv/bin/pytest tests/test_cost_tracking.py -q
docker stop ow-cost-tests
cd openweights/dashboard/frontend && npm ci && npm run build
```

Each database test creates and drops a private database in the disposable container.
It uses the repository's actual core schema/RLS and migration with a local `auth.uid`
stub; no configured Supabase credentials are used. Tests cover equal allocation,
retries, multiple workers, unallocated overhead, unknown prices, history retention,
fixed timestamps, tenant isolation, grants, admin permissions and budget enforcement.
17 changes: 17 additions & 0 deletions docs/release-0.13.0.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
# OpenWeights 0.13.0

- Track estimated USD costs by job, worker, submitting API key and user.
- Separate startup/idle overhead, allocating it equally across a worker's distinct
jobs in the default view, including failed and retried jobs.
- Configure lifetime per-key spending limits as an organization admin. Enforce
admission in PostgreSQL and cancel exhausted keys' work in the manager loop.
- Add a Costs dashboard with grouping, execution-only view, unknown-price warnings,
rate sources and limit management, plus `ow.costs` SDK and dashboard REST APIs.
- Retain cost history independently of operational record deletion. Record confirmed
worker termination and retry failed termination without prematurely ending costs.

**Migration required before upgrading the manager:** apply
`supabase/migrations/20260922000000_cost_tracking.sql`, then roll out v0.13.0.
Pre-upgrade usage remains unpriced. Limits are based on accrued estimates and can
overshoot during concurrent execution, shutdown or outages; they are not invoice caps.
See [cost tracking](cost-tracking.md) for accounting rules and rollout instructions.
2 changes: 2 additions & 0 deletions openweights/client/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
from typing import Any, Dict, Optional

from openweights.client.chat import AsyncChatCompletions, ChatCompletions
from openweights.client.costs import Costs
from openweights.client.decorators import supabase_retry
from openweights.client.events import Events
from openweights.client.files import (
Expand Down Expand Up @@ -193,6 +194,7 @@ def __init__(
self.jobs = Jobs(self)
self.runs = Runs(self)
self.events = Events(self)
self.costs = Costs(self)
self.async_chat = AsyncChatCompletions(self, deploy_kwargs=self.deploy_kwargs)
self.sync_chat = ChatCompletions(self, deploy_kwargs=self.deploy_kwargs)
self.chat = self.async_chat if use_async else self.sync_chat
Expand Down
49 changes: 49 additions & 0 deletions openweights/client/costs.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
"""Organization compute-cost estimates and lifetime API-key budgets (USD)."""

from decimal import Decimal

from openweights.client.decorators import supabase_retry


class Costs:
def __init__(self, ow_instance):
self._ow = ow_instance

@supabase_retry()
def report(self, limit=100, offset=0):
"""Return live costs by job, worker, API key and submitting user."""
return (
self._ow._supabase.rpc(
"get_cost_report",
{
"org_id": self._ow.organization_id,
"row_limit": limit,
"row_offset": offset,
},
)
.execute()
.data
)

@supabase_retry()
def set_limit(self, api_token_id: str, amount_usd=None):
"""Set a lifetime key budget; None removes it. Requires an admin user JWT.

Limits stop admission and cancel work after accrued spend reaches the
threshold. They are not prepaid reservations or exact invoice caps.
"""
if amount_usd is not None:
amount = Decimal(str(amount_usd))
if not amount.is_finite() or amount < 0:
raise ValueError(
"Spending limit must be a finite nonnegative USD amount"
)
amount_usd = str(amount)
self._ow._supabase.rpc(
"set_spending_limit",
{
"org_id": self._ow.organization_id,
"key_id": api_token_id,
"amount_usd": amount_usd,
},
).execute()
24 changes: 24 additions & 0 deletions openweights/cluster/costs.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
"""Snapshot the entire pod's hourly compute rate, without multiplying it twice."""

from decimal import Decimal, InvalidOperation

from openweights.cluster.start_runpod import GPU_COST_PER_HOUR


def worker_cost_fields(pod, gpu, count, started_at):
rate = pod.get("costPerHr")
source = "runpod"
try:
rate = Decimal(str(rate))
if not rate.is_finite() or rate < 0:
raise ValueError("Invalid pod rate")
except (InvalidOperation, ValueError):
estimate = GPU_COST_PER_HOUR.get(gpu)
rate = Decimal(str(estimate)) * count if estimate is not None else None
source = "hardware_estimate" if rate is not None else "unknown"
return {
"pod_id": pod["id"],
"hourly_cost_usd": str(rate) if rate is not None else None,
"cost_rate_source": source,
"billing_started_at": started_at,
}
27 changes: 26 additions & 1 deletion openweights/cluster/org_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
OpenWeights,
)
from openweights.client.decorators import supabase_retry
from openweights.cluster.costs import worker_cost_fields
from openweights.cluster.start_runpod import (
HARDWARE_REGISTRY,
is_spending_limit_error,
Expand Down Expand Up @@ -366,6 +367,7 @@ def clean_up_unresponsive_workers(self, workers):
runpod.terminate_pod(worker["pod_id"])
except Exception as e:
logger.error(f"Failed to terminate pod {worker['pod_id']}: {e}")
continue # Keep accruing cost and retry termination next cycle.

# 4) Finally, mark the worker as 'terminated' in the DB
self._ow._supabase.table("worker").update({"status": "terminated"}).eq(
Expand Down Expand Up @@ -569,6 +571,9 @@ def scale_workers(self, running_workers, pending_jobs):
).eq("id", worker_id).execute()

try:
billing_started_at = datetime.now(
timezone.utc
).isoformat()
pod = runpod_start_worker(
gpu=gpu,
count=count,
Expand All @@ -578,9 +583,22 @@ def scale_workers(self, running_workers, pending_jobs):
name=f"{self._ow.org_name}-{time.time()}-ow-1day",
runpod_client=runpod,
)
if pod.get("costPerHr") is None:
try:
details = runpod.get_pod(pod["id"])
pod["costPerHr"] = (details or {}).get(
"costPerHr"
)
except Exception as price_error:
logger.warning(
"Pod price unavailable; using hardware estimate: %s",
price_error,
)
self.hardware_registry.record_success(hardware_type)
self._ow._supabase.table("worker").update(
{"pod_id": pod["id"]}
worker_cost_fields(
pod, gpu, count, billing_started_at
)
).eq("id", worker_id).execute()
break
except Exception as e:
Expand Down Expand Up @@ -662,6 +680,12 @@ def set_shutdown_flags(self, idle_workers):
f"Failed to set shutdown flag for worker {idle_worker['id']}: {e}"
)

@supabase_retry()
def enforce_spending_limits(self):
return self._ow._supabase.rpc(
"enforce_spending_limits", {"org_id": self.org_id}
).execute()

def manage_cluster(self):
"""Main loop for managing the organization's cluster."""
logger.info(f"Starting cluster management for organization {self.org_id}")
Expand All @@ -674,6 +698,7 @@ def manage_cluster(self):
runpod.api_key = worker_env["RUNPOD_API_KEY"]
# try:
# Get active workers and pending jobs
self.enforce_spending_limits()
running_workers = self.get_running_workers()
pending_jobs = self.get_pending_jobs()

Expand Down
Loading
Loading