跳到主要内容

Quick Start

为什么是 Dex? 里同一条 order-processing Flow,从头跑一遍:安装 Dex,注册 Worker 和 Client,用 HTTP controller 启动 Flow,批准发货,再在 Dex Web 里看 Execution graph。

1. 安装 Dex

brew install superdurable/tap/dexcli
brew update && brew upgrade superdurable/tap/dexcli # when you want the latest
dexcli dev

dexcli 会启动 Dex Server、Dex Web,以及内部的 workflow backend。这些端口空闲时的默认地址:

ServiceAddress
Dex Webhttp://127.0.0.1:8802
Dex Server127.0.0.1:8801

同一台机器上再开一个 dexcli dev 会绑到别的空闲端口,并用自己的 SQLite。把打印出来的 Dex Server 地址传给 --server。flag 和持久化选项见 CLI README

2. Flow、Step、RPC、constructor 依赖

先扣款,再等卖家批准(Channel 或提醒 Timer),然后发货。发货 retry 用尽就退款。设计见 为什么是 Dex?

把 mock 的支付和发货 API 从 constructor 传进去。不要用 setter,也不要靠 global。

order_status = Attribute(
"order-status",
str,
AttributeIndex(IndexType.KEYWORD),
)
seller_ok = Channel[str]("seller-ok", str)


class ChargeStep(Step[OrderRequest]):
def __init__(self, service: MyDependencyService, ship: ShipStep) -> None:
self.service = service
self.ship = ship

def get_step_options(self) -> StepOptions:
return StepOptions(
execute_retry=RetryPolicy(
total_duration=timedelta(hours=1),
)
)

def execute(self, context: Context, input: OrderRequest) -> StepDecision:
self.service.charge_user(input.email, input.customer_id, input.amount)
order_status.set(context, "charged")
return go_to(ShipStep, input)


class ShipStep(Step[OrderRequest]):
def __init__(self, service: MyDependencyService, refund: RefundStep) -> None:
self.service = service
self.refund = refund

def get_step_options(self) -> StepOptions:
return StepOptions(
execute_retry=RetryPolicy(
total_duration=timedelta(hours=1),
)
).on_execute_failure_proceed_to(
RefundStep,
StepOptions(
execute_retry=RetryPolicy(
total_duration=timedelta(hours=1),
)
),
)

def wait_for(self, context: Context, input: OrderRequest) -> Wait:
return Wait.any_of(
seller_ok.for_one(),
Timer.by_duration(timedelta(hours=24)),
)

def execute(self, context: Context, input: OrderRequest) -> StepDecision:
if context.has_timer_fired():
self.service.send_email(
input.email,
"Reminder: approve shipment",
"Please approve or provide a tracking number.",
)
return go_to(ShipStep, input)
self.service.ship_item(input.order_id)
order_status.set(context, "shipped")
return graceful_complete(f"shipped:{input.order_id}")


class RefundStep(Step[OrderRequest]):
def __init__(self, service: MyDependencyService) -> None:
self.service = service

def execute(self, context: Context, input: OrderRequest) -> StepDecision:
self.service.update_external_system(f"refund {input.order_id}")
order_status.set(context, "refunded")
return graceful_complete(f"refunded:{input.order_id}")


class OrderProcessingFlow(Flow[OrderRequest]):
def __init__(self, service: MyDependencyService) -> None:
self.service = service
self.refund = RefundStep(service)
self.ship = ShipStep(service, self.refund)
self.charge = ChargeStep(service, self.ship)

def get_steps(self) -> StepList[OrderRequest]:
return StepList.start_step(self.charge).other_steps(self.ship, self.refund)

def get_persistence_schema(self) -> PersistenceSchema:
return PersistenceSchema.of(order_status, seller_ok)

@rpc
def approve(self, context: Context, _note: str) -> RPCResult[str]:
seller_ok.publish(context, "approved")
return RPCResult("ok")

@rpc
def describe(self, context: Context) -> RPCResult[str]:
return RPCResult(order_status.get(context))

例子: examples/python/dex_examples/products/order-processing

3. Worker 和 Client

把 Registry、Worker、Client 指到 dexcli127.0.0.1:8801)。用 service 构造 Flow,再注册。

        service = MyDependencyService()
pattern_service = ServiceDependency()

self.money_transfer = MoneyTransferFlow(service)
self.order_processing = OrderProcessingFlow(service)
        self.registry = Registry(tuple(flows), allow_async_handlers=True)
config.blob_cache_dir.mkdir(parents=True, exist_ok=True)
self.blob_cache = open_blob_cache(
BlobCacheConfig(str(config.blob_cache_dir), 1 << 30)
)
worker_options = WorkerOptions(
bind_address=config.worker_bind_address,
server_address=config.server_address,
worker_target=(
WorkerTarget(config.worker_target)
if config.worker_target
else None
),
)
self.worker = AsyncWorker(self.registry, self.blob_cache, worker_options)
self._client = AsyncClient(
self.registry,
self.blob_cache,
ClientOptions(
server_address=config.server_address,
worker_target=self.worker.worker_target,
),
)

例子: examples/python/dex_examples/products/order-processing

4. Controller

start handler 启动 Flow,再用 Client 的 wait-for-Step-completion API 等 ChargeStep。approve 往卖家 Channel 发消息。describe 读 order-status。Client 和 Flow 来自 constructor 或 factory 参数。

    @blueprint.get("/start")
async def start() -> Response:
flow_id = new_flow_id("order-processing")
request = OrderRequest(
flow_id,
"buyer@example.com",
"customer-1",
42,
)
run_id = await app_state.client.start_flow(
app_state.order_processing,
flow_id,
request,
start_options(),
)
await app_state.client.wait_for_step_completion(
flow_id,
CHARGE_STEP,
CHARGE_WAIT_TIMEOUT,
)
return started_flow(flow_id, run_id)

@blueprint.get("/approve")
async def approve() -> Response:
output = await app_state.client.invoke_rpc(
app_state.order_processing.approve,
required_query("workflowId"),
optional_query("notes", ""),
)
return jsonify(output)

@blueprint.get("/describe")
async def describe() -> Response:
flow_id = required_query("workflowId")
status = await app_state.client.invoke_rpc(
app_state.order_processing.describe,
flow_id,
)
return jsonify({"flowID": flow_id, "status": status})

例子: examples/python/dex_examples/products/order-processing

5. 运行

保持 dexcli dev 在跑。另开一个终端:

cd examples/python
uv run python main.py

然后启动一笔订单,打开 Dex Web,再批准发货:

curl -s 'http://127.0.0.1:8080/products/order-processing/start'

从 JSON 里复制 flowID。打开 http://127.0.0.1:8802,找到这条 Flow,打开详情。ShipStep 在等待。

curl -s "http://127.0.0.1:8080/products/order-processing/approve?workflowId=FLOW_ID"

FLOW_ID 换成 start 的返回值。这条 Flow 会发货并完成。

6. Dex Web

start 并 approve 之后,打开这条 Flow。Execution graph 里是 Charge 然后 Ship:

approve 之后的 Dex Web Execution graph

Timeline tab 列出同一次 run:

approve 之后的 Dex Web Timeline

接下来