| title | Running Pipelines |
|---|---|
| sidebar_position | 3 |
Start a pipeline, watch its progress, and stop it. Method tables live in the API reference; this page covers the workflow.
use() starts a pipeline from a file or an in-memory config and returns a dict whose
'token' identifies the running task — every data and control call takes it.
result = await client.use(filepath='pipeline.pipe')
token = result['token']Beyond filepath/pipeline, use() accepts token (a custom task token; the
server generates one when omitted), source, threads, use_existing,
args, ttl, pipelineTraceLevel (trace verbosity for the
run log), name (a display name for the task), and env
(per-run variable overrides). Pass the pipeline config as-is — the client sends it
to the server, which resolves ${ROCKETRIDE_*} variables from its merged
environment.
Running the same pipeline more than once at a time. Without a custom token,
the server names the task after its owner, project_id and source, so a second
use() of the same pipe fails with Pipeline is already running. Give each
instance its own token and they run side by side:
result = await client.use(filepath='pipeline.pipe', token=f'tk_{uuid.uuid4().hex}')The value you choose is the run's
private token: full control of
the run for anyone who presents it. So make it unguessable, and keep the tk_
prefix (without it the token cannot authenticate on webhook or dropper
endpoints). Tokens share one namespace across every user of the server: a
predictable value collides with other people's tasks, and with use_existing a
token that matches a running task attaches to it, whoever started it.
What two instances still share:
- The
pk_public authorization key. Both get the same one, and a webhook call or dropper upload that uses it always reaches the instance that started first. To send data to one specific instance, use that instance'stk_token. get_task_token(), keyed by the pipe'sproject_idand source, returns only one of them.- The monitor subscription. Both instances' events arrive in one
subscription. Each event carries the emitting run's id in
body.__id(the first eight characters of the token aftertk_, plus the source id); it is truncated and can collide, so treat it as a best-effort way to tell them apart. A few minutes after one instance ends, when the server removes it from its registry, the shared subscription is dropped and the surviving instance's events stop arriving; resubscribe. - The run log, written per identity and not built for two instances at once; treat a forked run's log as unreliable.
Check reused before trusting the result. use_existing returns the
instance that is already running under that token rather than starting the one
you submitted, and the result's reused flag is True when that happened. A
reused instance keeps the configuration it was created with — the pipeline in
this call is ignored, edits included — along with whatever state it has
accumulated. Benchmarks and A/B comparisons are where an unnoticed reuse costs
the most. Call restart() to apply new configuration to a live token.
Why a token: the server runs each pipeline as a separate task. The token targets
send(), send_files(), pipe(), chat(), get_task_status(), and terminate()
at the correct pipeline.
Poll get_task_status(token) — it returns completedCount, totalCount,
completed, state, exitCode, and more:
while True:
status = await client.get_task_status(token)
print(f'Progress: {status.get("completedCount", 0)}/{status.get("totalCount", 0)}')
if status.get('completed'):
break
await asyncio.sleep(2)For push-style progress instead of polling, add a monitor subscription; events
arrive at your on_event callback:
await client.add_monitor({'token': token}, ['apaevt_status_upload', 'apaevt_status_processing'])
# ... later:
await client.remove_monitor({'token': token}, ['apaevt_status_upload', 'apaevt_status_processing'])add_monitor(key, types) / remove_monitor(key, types) are reference-counted —
adding the same key merges types, removing unsubscribes a type only when its count
reaches zero. The key is {'token': ...} for a running task, or
{'project_id': ..., 'source': ...} (optionally with 'pipe_id' and/or
'team_id' — a team ID addresses that team's deployed run). The older
set_events(token, event_types, pipe_id=None) still works but is deprecated in
favor of the monitor pair.
validate(pipeline, source=None) checks a pipeline config server-side without
starting it and returns errors and warnings — cheap insurance before use().
terminate(token) stops the pipeline and frees server resources. Long-lived tasks
without a ttl run until terminated.
get_services() returns lightweight summaries of every service the server
supports (plus a deduplicated icon table and the server version). For a full
definition — config schema included — fetch one by name with get_service(name).
Note get_service raises on failure (ValueError for an empty name,
RuntimeError for an unknown service); it never returns None.
services = await client.get_services()
ocr = await client.get_service('ocr') # raises if unknownping() performs a liveness check against the server and raises on failure.
Deploying a pipeline so it persists server-side and runs on a schedule is a separate surface — see Deployments.