Common Patterns
Patterns observed across production scripts.
Downloading a file from blob storage
Input parameters that reference files arrive as either a raw blob storage URL or a JSON string of the form {"FileLocation": "https://..."}. Parse the bucket name and file ID from the URL and pass them to api.blob_storage:
import re
from uuid import UUID
async def _download_file(api: IAsyncScriptingApi, file_ref: str) -> bytes:
# Unwrap {"FileLocation": "..."} if needed
stripped = file_ref.strip()
if stripped.startswith("{"):
import json
stripped = json.loads(stripped)["FileLocation"]
# Extract bucket name and file ID from the URL path
m = re.search(r"/buckets/([^/]+)/files/([^/?]+)", stripped)
if not m:
raise ValueError(f"Cannot parse file reference: {file_ref!r}")
result = await api.blob_storage.bucket_files_v1.bucket_download_file(
bucket_name=m.group(1),
file_id=UUID(m.group(2)),
)
return result.contentUploading a file to the workflow order
Upload the file, then immediately mark it PROCESSED so downstream steps can consume it:
from uuid import UUID
from cellario_cloud_data.models import FileUploadRequest, FileUpdateRequest, FileStatus
async def _upload_file(
api: IAsyncScriptingApi,
filename: str,
content: bytes,
media_type: str = "application/octet-stream",
) -> None:
order_id = UUID(api.workflow_order_id)
result = await api.data.v2_workflow_orders.upload_file(
order_id=order_id,
body=FileUploadRequest(file_name=filename, content=content, content_type=media_type),
)
await api.data.v2_workflow_orders.update_file(
order_id=order_id,
file_id=result.id,
body=FileUpdateRequest(sort_order=0, status=FileStatus.PROCESSED),
)Reading secrets
Secrets are stored in the Cellario vault and injected at runtime. Access them via api.secrets. A missing key returns "" rather than raising.
raw = api.secrets["my-service-api-key"]
if not raw:
raise ValueError("Secret 'my-service-api-key' is not configured")If the secret is a GCP service account JSON (or any multi-line PEM), normalize escaped newlines after reading from the vault:
import json
sa = json.loads(api.secrets["gcp-service-account"])
sa["private_key"] = sa["private_key"].replace("\\n", "\n")Publishing an event
import json
from datetime import datetime, timezone
from cellario_cloud_events.models import (
EventsApiEventsCreateRequest,
EventsApiEventsEventSourceCategory,
EventsApiEventsEventSourceComponent,
EventsApiEventsEventSeverity,
)
body = EventsApiEventsCreateRequest(
type_="my.step.completed",
source_name="my-script",
source_id=api.workflow_order_id,
source_created_at=datetime.now(timezone.utc),
source_category=EventsApiEventsEventSourceCategory.CELLARIO_OS,
source_component=EventsApiEventsEventSourceComponent.OTHER,
source_component_extended_description="my-script-name",
stream_id=api.workflow_order_id,
payload_type="my.step.completed.v1",
payload=json.dumps({"order_id": api.workflow_order_id}),
severity=EventsApiEventsEventSeverity.INFO,
)
await api.events.events.create(body=body)Passing collections between steps
list and dict outputs are silently dropped by the platform. Serialize to a JSON string instead, using a _json suffix by convention:
@dataclass
class Result:
plates_json: Annotated[
str,
OutParam(display_name="Plates JSON", description="JSON-encoded list of plate records"),
]
# Serialise on the way out
return Result(plates_json=json.dumps(plates))
# Deserialise on the way in (next step)
plates = json.loads(plates_json)For workflows with parallel branches that rejoin, use the merge step pattern — each branch emits its own _json output, the merge step takes both as inputs and combines them:
@Main
async def merge_branch_records(
api: IAsyncScriptingApi,
branch_a_json: Annotated[str, InParam(display_name="Branch A", description="Records from branch A")] = "[]",
branch_b_json: Annotated[str, InParam(display_name="Branch B", description="Records from branch B")] = "[]",
) -> Result:
records = json.loads(branch_a_json) + json.loads(branch_b_json)
return Result(merged_json=json.dumps(records))Loop control
Wire a bool output named keep_going to the workflow's loop condition node:
@dataclass
class Result:
result_json: Annotated[str, OutParam(display_name="Result", description="Processed batch")]
keep_going: Annotated[bool, OutParam(display_name="Keep Going", description="True while items remain")]
status_message: Annotated[str, OutParam(display_name="Status", description="Progress summary")]
@Main
async def process_batch(
api: IAsyncScriptingApi,
items_json: Annotated[str, InParam(display_name="Items", description="Remaining items as JSON")] = "[]",
) -> Result:
items = json.loads(items_json)
current, remaining = items[0], items[1:]
result = await _process(api, current)
return Result(
result_json=json.dumps(result),
keep_going=len(remaining) > 0,
status_message=f"Processed '{current}'. {len(remaining)} item(s) remaining.",
)Retrying external HTTP calls
Use tenacity for outbound HTTP. Retry only on transient failures (timeouts and 5xx); let 4xx errors fail immediately:
import httpx
from tenacity import retry, retry_if_exception, stop_after_attempt, wait_exponential
_with_retry = retry(
retry=retry_if_exception(
lambda e: (
isinstance(e, httpx.TimeoutException)
or (isinstance(e, httpx.HTTPStatusError) and e.response.status_code >= 500)
)
),
stop=stop_after_attempt(3),
wait=wait_exponential(multiplier=1, min=1, max=10),
reraise=True,
)
@_with_retry
async def _call_external(client: httpx.AsyncClient, url: str) -> dict:
r = await client.get(url)
r.raise_for_status()
return r.json()Collecting multiple validation errors
Rather than raising on the first error, accumulate all failures and surface them together:
errors: list[str] = []
for plate_id in plate_ids:
plate = lookup_plate(plate_id)
if plate is None:
errors.append(f"No plate found with plateId {plate_id!r}")
elif plate.status != "Active":
errors.append(f"Plate {plate_id!r} has status {plate.status!r}, expected 'Active'")
if errors:
bullet_list = "\n".join(f" - {e}" for e in errors)
raise ValueError(f"{len(errors)} validation error(s):\n{bullet_list}")Writing testable scripts
Extract all logic into a pure synchronous (or async) evaluate() function. The @Main entry point becomes a thin wrapper that feeds it api context:
def evaluate(plates_json: str, threshold: int) -> Result:
plates = json.loads(plates_json)
valid = [p for p in plates if p["count"] >= threshold]
return Result(valid_json=json.dumps(valid), count=len(valid))
@Main
async def filter_plates(
api: IAsyncScriptingApi,
plates_json: Annotated[str, InParam(display_name="Plates", description="Plates as JSON")] = "[]",
threshold: Annotated[int, InParam(display_name="Threshold", description="Minimum count")] = 1,
) -> Result:
return evaluate(plates_json, threshold)Tests then call evaluate() directly with no mock needed:
import asyncio
from unittest.mock import AsyncMock
def test_filters_below_threshold():
plates = json.dumps([{"id": "A", "count": 5}, {"id": "B", "count": 1}])
result = evaluate(plates, threshold=3)
assert result.count == 1
# For scripts that must be tested end-to-end with api:
def test_with_api():
api = AsyncMock()
result = asyncio.run(filter_plates(api, plates_json='[{"id":"A","count":5}]', threshold=3))
assert result.count == 1Mock mode
For scripts that call external systems that may not be available in all environments, add a mock_mode: bool input. The platform sends False when unconfigured; set = True as the signature default if you want it to default to mock in local testing. Build the full payload regardless — this validates the logic — but skip the outbound call:
@Main
async def submit_to_external(
api: IAsyncScriptingApi,
data_json: Annotated[str, InParam(display_name="Data", description="Payload to submit")] = "{}",
mock_mode: Annotated[bool, InParam(display_name="Mock Mode", description="Skip the external call")] = True,
) -> Result:
payload = _build_payload(json.loads(data_json))
if mock_mode:
return Result(submission_id=f"MOCK-{api.workflow_order_id}", success=True)
submission_id = await _post_to_external(payload)
return Result(submission_id=submission_id, success=True)