Skip to content

Build and run a flow#

This page covers writing a flow three ways — the console's editor, YAML with astra, Python — and running and following its executions. The concepts are in Flows; every field is in Flow specification.

Before you begin#

  • The editor or admin role (viewers watch executions).
  • The agents, flows and drives the flow names, in the workspace.
  • For the CLI: astra login and astra use <org>/<workspace> (see Install the CLI).
  • For the API: ASTRA_TOKEN and WS as in the REST API page.

Write the flow#

This flow collects a day's incidents, has an agent summarise them, asks an admin, then publishes:

nightly-digest.yaml
apiVersion: anemoi.astralyx/v1
kind: Flow
metadata:
  name: nightly-digest
spec:
  description: Collect yesterday's incidents, summarise them, ask, then publish.
  drive: digests                      # a shared drive: every step sees /flow
  inputs:
    - name: day
      description: The day to report on (YYYY-MM-DD)
      required: true
  timeout_seconds: 21600
  steps:
    - id: collect
      kind: run
      retry: {max: 3, backoff_seconds: 60}
      timeout_seconds: 1800
      template:
        image: ghcr.io/acme/incident-export:2.1
        command: export
        args: ["--day", "${{ inputs.day }}", "--out", "/flow/collect/incidents.json"]
        requested_resources: {cpu_cores: 1, memory_bytes: 268435456}
    - id: summarise
      kind: agent
      needs: [collect]
      when: steps.collect.outputs.count != '0'
      agent: researcher
      input: "Summarise the ${{ steps.collect.outputs.count }} incidents of ${{ inputs.day }} in /flow/collect/incidents.json."
      retry: {max: 2, backoff_seconds: 120}
    - id: ok
      kind: approval
      needs: [summarise]
      when: steps.summarise.state == 'Succeeded'
      target: "Publish the incident digest for ${{ inputs.day }}?"
      approvers: ["role:admin"]
      wait_seconds: 28800
    - id: publish
      kind: run
      needs: [ok]
      when: steps.summarise.state == 'Succeeded'
      template:
        image: ghcr.io/acme/digest-publisher:1.0
        command: publish
        args: ["/flow/summarise/answer.md"]
        requested_resources: {cpu_cores: 1, memory_bytes: 268435456}
  outputs:
    incidents: "${{ steps.collect.outputs.count }}"

The collect step writes its outputs with echo "count=3" >> "$ASTRAEUS_OUTPUTS"; the agent's answer lands in /flow/summarise/answer.md.

  1. Open Anemoi → Flows and press New flow.
  2. Name: nightly-digest.
  3. On Steps, Add each step — Agent, Run, Approval, Wait for event, Sleep, Subflow — and fill it in on the right (Id, Description, Kind and its fields, needs, when, retry, repeat, timeout). Click the dot on a step's right, then another step, to make the second need the first; click an arrow to remove it.
  4. On Inputs & outputs, declare day and the output incidents.
  5. YAML shows the same flow as a file; you can paste a flow there.
  6. Problems are listed under the editor as you go, checked by the cluster too. Press Create the flow.

The flow editor: steps as a graph, the selected step's fields on the right

The same flow on the YAML tab

$ astra anemoi flows validate -f nightly-digest.yaml
nightly-digest.yaml: flow nightly-digest is valid
$ astra anemoi flows validate -f nightly-digest.yaml --remote
nightly-digest.yaml: flow nightly-digest is valid
$ astra anemoi flows apply -f nightly-digest.yaml
flow nightly-digest: version 1 written, now current

validate runs the cluster's own checks on the file, here; --remote also asks the cluster, which knows whether the agents, flows and drives named exist. apply creates the flow, or writes its next version (--no-current keeps the current one). A flow file may also be JSON.

The astralyx package (Python 3.9 or newer, standard library only):

nightly_digest.py
from astralyx.anemoi import Client, Flow, Retry

flow = Flow("nightly-digest", description="Collect, summarise, ask, publish.",
            drive="digests", timeout_seconds=21600)
day = flow.input("day", description="The day to report on (YYYY-MM-DD)", required=True)

collect = flow.run("collect", retry=Retry(max=3, backoff_seconds=60), timeout_seconds=1800, template={
    "image": "ghcr.io/acme/incident-export:2.1", "command": "export",
    "args": ["--day", day, "--out", "/flow/collect/incidents.json"],
    "requested_resources": {"cpu_cores": 1, "memory_bytes": 256 << 20}})
summarise = flow.agent("summarise", agent="researcher", needs=[collect],
                       when=f"{collect.outputs.count} != '0'", retry=Retry(max=2, backoff_seconds=120),
                       input=f"Summarise the {collect.output('count')} incidents of {day} in /flow/collect/incidents.json.")
ok = flow.approval("ok", needs=[summarise], when=f"{summarise.state} == 'Succeeded'",
                   target=f"Publish the incident digest for {day}?", approvers=["role:admin"], wait_seconds=28800)
flow.run("publish", needs=[ok], when=f"{summarise.state} == 'Succeeded'", template={
    "image": "ghcr.io/acme/digest-publisher:1.0", "command": "publish",
    "args": ["/flow/summarise/answer.md"],
    "requested_resources": {"cpu_cores": 1, "memory_bytes": 256 << 20}})
flow.output("incidents", collect.output("count"))

assert flow.validate() == [], flow.validate()
c = Client.from_config()           # what `astra login` and `astra use` saved
c.flows.apply(flow)                 # a new flow, or its next version
print(c.flows.check(flow))          # the cluster's checks: {ok, errors, warnings}

step.output(key) is ${{ steps.<id>.outputs.<key> }}, for templates; step.outputs.<key> and step.state are the bare names, for when and until. flow.to_json() is the same document as the YAML file. Client.from_config() reads ~/.config/astra/config.json and the ASTRA_URL, ASTRA_TOKEN, ASTRA_ORG, ASTRA_WORKSPACE and ASTRA_CLUSTER variables.

$ yq -o=json '{"metadata": .metadata, "spec": .spec}' nightly-digest.yaml \
    | curl -sS -X POST "$WS/flows" -H "Authorization: Bearer $ASTRA_TOKEN" \
        -H "Content-Type: application/json" -d @-

POST $WS/flow-checks with {"name", "spec"} checks without writing: {"ok", "errors": [{"path", "message"}], "warnings": [...]}. A new version is PUT $WS/flows/{name} with {"spec", "current"}.

Without a drive, the check warns: each execution gets a drive of its own, a copy on each machine: steps on different machines do not see each other's files. Name a shared drive to share them.

Run it#

On the flow's page press Run flow, fill in its inputs (and a Version other than the current one if you like) and confirm. The execution's page opens.

$ astra anemoi flows run nightly-digest --input day=2026-09-30 --wait
execution nightly-digest-cb4855 of flow nightly-digest started
nightly-digest-cb4855: Running
nightly-digest-cb4855: Waiting

--input key=value (repeatable), --version N, --wait to follow it until it ends and print its steps (exit status 1 unless it Succeeded).

e = c.flows.run("nightly-digest", inputs={"day": "2026-09-30"})
done = c.executions.wait(e["metadata"]["name"], timeout=6 * 3600)
print(done["status"]["state"], done.get("outputs"))
$ curl -sS -X POST "$WS/flows/nightly-digest/executions" -H "Authorization: Bearer $ASTRA_TOKEN" \
    -H "Content-Type: application/json" -d '{"inputs": {"day": "2026-09-30"}}' | jq -r .metadata.name
nightly-digest-cb4855

Follow an execution#

The execution's page draws the graph live: each step's state (pending, running, waiting, succeeded, failed, skipped), its attempts with links to their runs and approvals, the machine, the outputs; then the execution's Inputs, Outputs, the Receipts of its agent runs and the Events it received. Click a step to see only it.

An execution: the graph, each step's attempts, its inputs and outputs

$ astra anemoi flows executions --flow nightly-digest
$ astra anemoi flows get nightly-digest-cb4855
$ astra anemoi flows get nightly-digest-cb4855 --json

get prints the state, then STEP STATE ATTEMPTS LAST MACHINE WHY OUTPUTS, the receipts of its agent runs, and its outputs.

$ curl -sS "$WS/flow-executions?flow=nightly-digest&state=Running" -H "Authorization: Bearer $ASTRA_TOKEN"
$ curl -sS "$WS/flow-executions/nightly-digest-cb4855" -H "Authorization: Bearer $ASTRA_TOKEN" \
    | jq '{state: .status.state, steps: (.steps | map_values({state, reason, attempts: (.attempts | length)}))}'

Send an event, retry a step, cancel#

On the execution's page: Send event on a waiting wait_event step (its name and key=value data); Retry step on a failed step (a new attempt); Cancel stops its running steps.

$ astra anemoi flows events send nightly-digest-cb4855 deployed env=prod sha=4be1c2
execution nightly-digest-cb4855: event deployed sent
$ astra anemoi flows retry nightly-digest-cb4855 summarise
execution nightly-digest-cb4855: step summarise runs again
$ astra anemoi flows cancel nightly-digest-cb4855
execution nightly-digest-cb4855: cancelling
c.executions.send_event("nightly-digest-cb4855", "deployed", {"env": "prod"})
c.executions.retry("nightly-digest-cb4855", "summarise")
c.executions.cancel("nightly-digest-cb4855")
$ curl -sS -X POST "$WS/flow-executions/nightly-digest-cb4855/events/deployed" -H "Authorization: Bearer $ASTRA_TOKEN" \
    -H "Content-Type: application/json" -d '{"data": {"env": "prod"}}'
$ curl -sS -X POST "$WS/flow-executions/nightly-digest-cb4855/steps/summarise/retry" -H "Authorization: Bearer $ASTRA_TOKEN"
$ curl -sS -X POST "$WS/flow-executions/nightly-digest-cb4855/cancel" -H "Authorization: Bearer $ASTRA_TOKEN"

An event's data is key=value text, at most 4 KiB. A retried approval or event step asks again. DELETE $WS/flow-executions/{name} removes an execution that ended, with its own drive.

Change a flow#

Edit (new version) on the flow's page opens the editor with the current version; tick Runs use this version from now on or not, and press Write the version. Versions lists them, with Make current.

astra anemoi flows apply -f again writes the next version; astra anemoi flows list shows FLOW CURRENT LATEST AGE.

$ curl -sS -X PUT "$WS/flows/nightly-digest/current" -H "Authorization: Bearer $ASTRA_TOKEN" \
    -H "Content-Type: application/json" -d '{"version": 1}'

Executions keep the version they started with.

Run a flow on a schedule#

Flows have no schedule of their own yet. Start executions from outside — your CI, a cron job — with astra anemoi flows run and a personal API token in ASTRA_TOKEN, or the API. See A nightly flow with retries.