DEV Community

Cover image for An AWS Step Functions Compiler: Python In, Readable JSON Out

An AWS Step Functions Compiler: Python In, Readable JSON Out

Changing one retry policy in a Step Functions workflow can mean editing the same block of JSON on every Task that uses it. A workflow is defined in the Amazon States Language (ASL), a JSON document of states that name each other, with the logic in expressions. Writing it by hand means wiring every Next and repeating the same error handling on every Task.

sfnx lets you write the workflow as a Python function instead, and compiles it to a definition in JSONata mode that you deploy with your usual tools. The output is meant to read like a definition a person would write: few states, named after what they do.

This post uses sfnx v3.0.0.

The 25-Copy Retry Rule

In one real workflow, the same retry policy appears on 25 Tasks. It is populate_draft_data_sfn_template.asl.json from OrcaBus, an open-source platform for genomics pipelines: a JSONata ASL definition of 1,576 lines and 104 states.

All 25 are Lambda Tasks. 22 of the policies are exact copies, and three add States.TaskFailed to the list:

"Retry": [
  {
    "ErrorEquals": [
      "Lambda.ServiceException",
      "Lambda.AWSLambdaException",
      "Lambda.SdkClientException",
      "Lambda.TooManyRequestsException"
    ],
    "IntervalSeconds": 1,
    "MaxAttempts": 3,
    "BackoffRate": 2,
    "JitterStrategy": "FULL"
  }
],
Enter fullscreen mode Exit fullscreen mode

Those 25 blocks are about a fifth of the file. ASL itself has no way to define a retrier once and refer to it, so changing the backoff in this file means editing all 25 copies, and a copy that drifts is easy to miss.

One Rule in Python

In sfnx, a value assigned at the top of the module is written into the definition wherever it is read. The retry policy becomes one constant, and each Task names it. This is a small illustration of the pattern, not a port of the OrcaBus workflow:

from sfnx import aws, state_machine

LAMBDA_RETRY = [
    {
        "ErrorEquals": [
            aws.optimized.lambda_.errors.ServiceException,
            aws.optimized.lambda_.errors.AWSLambdaException,
            aws.optimized.lambda_.errors.SdkClientException,
            aws.optimized.lambda_.errors.TooManyRequestsException,
        ],
        "IntervalSeconds": 1,
        "MaxAttempts": 3,
        "BackoffRate": 2,
        "JitterStrategy": "FULL",
    }
]


@state_machine
def populate(input: dict):
    checked = aws.optimized.lambda_.invoke(
        FunctionName="validate-draft", Payload=input, retry=LAMBDA_RETRY
    )
    tags = aws.optimized.lambda_.invoke(
        FunctionName="get-library-tags", Payload=input, retry=LAMBDA_RETRY
    )
    return {"valid": checked["Payload"]["isValid"], "tags": tags["Payload"]}
Enter fullscreen mode Exit fullscreen mode

The compiled definition still has the full Retry on both Tasks, because that is what Step Functions reads. What changes is where you edit it: once, in the source.

An Order in Five States

Here is a whole workflow. It reserves every item of an order in DynamoDB, fails with OutOfStock if an item runs out, and then charges for the order through Lambda. It is examples/orders.py in the repository. It is kept small to show the control flow: if a later item or the charge fails, the stock already reserved is not put back, which a real order workflow would need to handle.

from typing import TypedDict

from sfnx import Timeout, aws, state_machine


class Item(TypedDict):
    sku: str
    quantity: int


class Order(TypedDict):
    id: str
    items: list[Item]


class OutOfStock(Exception):
    pass


@state_machine(timeout=300)
def fulfill(input: Order):
    """Reserve every item of an order, then charge for it."""
    items = input["items"]
    for item in items:
        try:
            aws.sdk.dynamodb.update_item(
                TableName="stock",
                Key={"sku": {"S": item["sku"]}},
                UpdateExpression="SET quantity = quantity - :n",
                ConditionExpression="quantity >= :n",
                ExpressionAttributeValues={":n": {"N": str(item["quantity"])}},
                retry=[{"ErrorEquals": [Timeout], "MaxAttempts": 3}],
            )
        except aws.sdk.dynamodb.errors.ConditionalCheckFailedException:
            raise OutOfStock(f"{item['sku']} is out of stock") from None
    receipt = aws.optimized.lambda_.invoke(FunctionName="charge", Payload=input)
    return {"order": input["id"], "receipt": receipt["Payload"]}
Enter fullscreen mode Exit fullscreen mode

Save it as app.py and compile it. You need uv; uvx runs sfnx without installing it into a project:

uvx sfnx@3.0.0 compile app.py -o fulfill.asl.json
Enter fullscreen mode Exit fullscreen mode

The definition has five states:

State Type From the source
items Pass items = input["items"], and the loop counter starting at 0
for Choice the for loop: another item, or move on
updateItem Task the DynamoDB call, with its Retry, the except as its Catch, and the counter moving on in its Assign
raise Fail raise OutOfStock(...), with the SKU in the cause
receipt Task the Lambda call, whose Output is the return

An excerpt, with the arguments left out and some lines joined, shows the loop and the DynamoDB call:

"for": {
  "Type": "Choice",
  "Choices": [
    { "Condition": "{% $item_index < $count($items) %}", "Next": "updateItem" }
  ],
  "Default": "receipt"
},
"updateItem": {
  "Type": "Task",
  "Resource": "arn:aws:states:::aws-sdk:dynamodb:updateItem",
  "Arguments": { ... },
  "Retry": [{ "ErrorEquals": ["States.Timeout"], "MaxAttempts": 3 }],
  "Catch": [
    { "ErrorEquals": ["DynamoDb.ConditionalCheckFailedException"], "Next": "raise" }
  ],
  "Assign": { "item_index": "{% $item_index + 1 %}" },
  "Next": "for"
}
Enter fullscreen mode Exit fullscreen mode

The whole definition
{
  "Comment": "Reserve every item of an order, then charge for it.",
  "QueryLanguage": "JSONata",
  "TimeoutSeconds": 300,
  "StartAt": "items",
  "States": {
    "items": {
      "Type": "Pass",
      "Assign": {
        "items": "{% $states.context.Execution.Input.items %}",
        "item_index": 0
      },
      "Next": "for"
    },
    "for": {
      "Type": "Choice",
      "Choices": [
        {
          "Condition": "{% $item_index < $count($items) %}",
          "Next": "updateItem"
        }
      ],
      "Default": "receipt"
    },
    "updateItem": {
      "Type": "Task",
      "Resource": "arn:aws:states:::aws-sdk:dynamodb:updateItem",
      "Arguments": {
        "TableName": "stock",
        "Key": {
          "sku": {
            "S": "{% $items[$item_index].sku %}"
          }
        },
        "UpdateExpression": "SET quantity = quantity - :n",
        "ConditionExpression": "quantity >= :n",
        "ExpressionAttributeValues": {
          ":n": {
            "N": "{% $string($items[$item_index].quantity) %}"
          }
        }
      },
      "Retry": [
        {
          "ErrorEquals": [
            "States.Timeout"
          ],
          "MaxAttempts": 3
        }
      ],
      "Catch": [
        {
          "ErrorEquals": [
            "DynamoDb.ConditionalCheckFailedException"
          ],
          "Next": "raise"
        }
      ],
      "Assign": {
        "item_index": "{% $item_index + 1 %}"
      },
      "Next": "for"
    },
    "raise": {
      "Type": "Fail",
      "Error": "OutOfStock",
      "Cause": "{% $items[$item_index].sku & ' is out of stock' %}"
    },
    "receipt": {
      "Type": "Task",
      "Resource": "arn:aws:states:::lambda:invoke",
      "Arguments": {
        "FunctionName": "charge",
        "Payload": "{% $states.context.Execution.Input %}"
      },
      "Output": {
        "order": "{% $states.context.Execution.Input.id %}",
        "receipt": "{% $states.result.Payload %}"
      },
      "End": true
    }
  }
}
Enter fullscreen mode Exit fullscreen mode

In this example, the compiler folds the counter update into the Task's Assign and the return into the last Task's Output, avoiding separate states for those steps.

The repository also compares a larger workflow with its hand-written original: AWS's coding-agent sample, with 8 states, and its Python version, which compiles to 7. 12 mocked scenarios run both in the local runner, and in each of those 12 runs the two make the same calls, end the same way, and enter the same number of states. These are local runs, not a measurement of your bill: a Standard workflow is billed per state transition, retries included, so check the paths your own executions take.

States Named After Your Code

Each state is named after what it does: the variable it assigns (items, receipt), if, for, return or raise, or the API a Task calls on its own line (updateItem). A Lambda invoke, a DynamoDB putItem or an SNS publish whose result nothing assigns also gets its target, such as invoke charge or putItem orders. So the console graph and the execution history read like the source.

Comments go into the definition too. The function's docstring is the state machine's Comment, and comment lines right above a statement become the Comment of the first state it makes.

When the generated name is not the one you want, name the state yourself with a comment at the end of the line:

if answer["decision"] == "approve":  # state: IsApproved
    return "approved"
Enter fullscreen mode Exit fullscreen mode

A named state is also kept as its own state: the compiler does not merge it into another one.

Catch SDK Typos Early

The service, operation, argument and error names of AWS SDK integrations are checked against the botocore service models when you compile. A typo stops the compilation with the line, the column, and the name you probably meant. On separate compile runs, for example:

app.py:27:17: updateItem has no argument Tablename; did you mean TableName?
app.py:26:13: dynamodb has no operation updateitem; did you mean update_item or update_table or put_item?
app.py:34:16: dynamodb has no error ConditionalCheckFailed; did you mean ConditionalCheckFailedException?
Enter fullscreen mode Exit fullscreen mode

In hand-written ASL, an SDK error name can be misspelled yet remain valid ASL, leaving the Catch unable to match. The botocore checks also cannot confirm that Step Functions supports a modeled action.

The compiler also rejects Python it cannot translate, and says what to write instead:

app.py:6:9: calling print() is not supported; write it with operators or jsonata(), or compute it in a Lambda task
Enter fullscreen mode Exit fullscreen mode

Test Before Deploying

sfnx.testing runs a JSONata definition on your machine, with each Task answered by a function in your test. A test checks which calls the workflow makes, where it goes when a call fails, and what it returns, without AWS credentials:

from sfnx.compiler import compile_file
from sfnx.testing import Call, Failure, run

(ORDERS,) = compile_file("app.py").values()
ORDER = {"id": "o1", "items": [{"sku": "a", "quantity": 2}, {"sku": "b", "quantity": 1}]}
UPDATE = "arn:aws:states:::aws-sdk:dynamodb:updateItem"
INVOKE = "arn:aws:states:::lambda:invoke"


def stock(*out: str):
    """DynamoDB with the SKUs given out of stock, and a charge of 30."""

    def tasks(call: Call) -> object:
        assert isinstance(call.arguments, dict)
        if call.resource == UPDATE:
            if call.arguments["Key"]["sku"]["S"] in out:
                raise Failure("DynamoDb.ConditionalCheckFailedException", "declined")
            return {}
        assert call.resource == INVOKE
        return {"Payload": {"total": 30}}

    return tasks


def test_an_order_in_stock_is_charged():
    execution = run(ORDERS, ORDER, stock())
    assert execution.output == {"order": "o1", "receipt": {"total": 30}}
    assert [call.resource for call in execution.calls] == [UPDATE, UPDATE, INVOKE]


def test_an_item_out_of_stock_fails_the_order_before_charging():
    execution = run(ORDERS, ORDER, stock("b"))
    assert (execution.error, execution.cause) == ("OutOfStock", "b is out of stock")
    assert INVOKE not in [call.resource for call in execution.calls]
Enter fullscreen mode Exit fullscreen mode

Save it as test_orders.py in a tests directory next to app.py, then add sfnx and pytest to the project and run it from there:

uv init
uv add --dev "sfnx==3.0.0" pytest
uv run pytest
Enter fullscreen mode Exit fullscreen mode

The runner also works on JSONata definitions you did not write with sfnx. My previous post walks through testing an existing definition this way.

Where sfnx Stops

  • It stops at the definition. sfnx does not deploy and does not run workflows in AWS. You deploy the JSON with CDK, SAM, CloudFormation or the AWS CLI; the deployment guide has the snippets and the IAM actions each kind of Task needs.
  • It compiles a subset of Python. Control flow becomes states and expressions become JSONata, but arbitrary libraries do not compile; put that work in a Lambda function, or write JSONata yourself with jsonata(). Some results differ from what CPython would compute, and the language reference lists what compiles and where results differ.
  • A local run is not AWS. Time does not pass, Parallel branches and Map iterations run one at a time, and the JSONata engine is a Python one. The testing guide lists the differences.

Try It on a Workflow You Know

Start with the example in this post. Save orders.py and compile it:

uvx sfnx@3.0.0 compile orders.py
Enter fullscreen mode Exit fullscreen mode

Read the definition it prints, change a line, compile again, and see which states move.

Then pick a small workflow you already run, write it as a Python function, and compare the compiled definition with the one you have. If your existing definition uses JSONata, the runner can check that both make the same calls on the paths you care about.

sfnx is on GitHub. If you try it, I'd like to hear where the generated definition reads differently from yours.

Top comments (0)