RPC
RPC 是 “Remote Procedure Call” 的缩写,允许外部系统与 Flow execution 交互。
Client 发起 RPC。Dex Server 将它发送给 Worker。Worker 执行后,将结果返回给 Dex Server 持久化。Dex Server 再将输出返回给 Client。
定义和调用 RPC
RPC 可以通过 Attribute 读取和写入 durable state、向 Channel publish,也可以 trigger 一个全新的 Step 来运行。
每次 RPC 调用都是一个 commit 边界。Dex Server 会暂存 RPC 中的 Attribute 写入和 Channel publish,并只在 method 成功返回时,将它们与 RPC result 一起 commit。如果 RPC 失败,这些 durable 变更都不会被 commit。外部 API 调用不在这个原子边界内。
@rpc
def trigger(self, context: Context, input: str) -> RPCResult[str]:
self.data.set(context, input)
self.example_ch.publish(context, None)
return RPCResult(input, next_steps=(StepMovement.of(ExampleStep, input),))
Java 和 Kotlin
Client.newRpcStub 会创建 Flow 的子类。Java Flow class 和它的 RPC method 都不能是 final。Kotlin 默认将两者都设为 final,所以要将 Flow class 和它的 RPC method 声明为 open。
TypeScript codec
TypeScript 未设置 inputCodec 或 outputCodec 时会使用 JSON。input 或 output 不是 JSON 时需要提供 codec,包括 string、number、boolean 和 bigint 等 scalar value。
然后像下面的例子一样调用 RPC:
result = await client.invoke_rpc(flow.trigger, flow_id, message)
选择性加载 RPC 状态
每个 RPC 都会收到普通 Attribute value 和 Channel size metadata。读取 AttributeMap entry 或 pending Channel message 时必须显式加载,避免不使用这些大型 collection 的 RPC 仍要接收它们。
Handler 需要所有当前 instance 时,加载整个 AttributeMap 或 ChannelMap;只需要已知 key 时,加载指定 map instance。需要 pending message 的 ID 或 Value 时,加载对应 Channel。Channel size、ChannelMap key 和 size 无需加载 pending message 即可读取。
Python、Java、TypeScript 和 Rust 在 RPC definition 上声明 load;Go 在 InvokeOptions 中提供。Dex Server 会根据 Flow persistence schema 验证 requested definition,并随 RPC snapshot 返回已加载范围。Handler 读取选择性加载的 collection 时,SDK 会检查这个范围。
如果 AttributeMap entry 或 instance 不在该范围内,SDK 会抛出 AttributeMapNotLoadedError。如果 pending message 不在该范围内,SDK 会抛出 ChannelMessagesNotLoadedError。Java 使用 AttributeMapNotLoadedException 和 ChannelMessagesNotLoadedException。
显式加载的 collection 可以为空,这不是错误。空的 pending-message snapshot 会返回空列表。不存在的 AttributeMap entry 会使用 SDK 正常的缺失值语义。Load 会建立一份 input snapshot。它不会消费 message,也不会使 RPC 具备 transactional 语义或阻止并发 operation。
Transactional read 与 write
Transactional RPC 会把 RPC 的 read 和 write 作为一个 atomic operation 执行。
持有 Attribute lock 的 RPC 已经隐式具备 transactional 语义。也可以设置 is_transactional=true,显式启用 transactional RPC。这对 Channel message deletion 尤其有用,因为 Dex 会在 commit 任何 write 前验证 message 是否仍处于 pending 状态。
Caller 只传 pending message ID。RPC 加载 source Channel,从 snapshot 读取原始 Value,再暂存 deletion 和 destination publication。如果 concurrent consumer 已取走该消息,deletion validation 会失败,因此 Dex 不会 commit 任何一项 effect:
@rpc(is_transactional=True, load_channels=(queued_messages,))
def move_queued_message_to_prioritized_messages(
self, context: Context, queued_message: QueuedMessageReference
) -> None:
message_to_prioritize = self.queued_messages.find_pending_message(
context, queued_message.message_id
)
self.queued_messages.delete(context, queued_message.message_id)
if message_to_prioritize is not None:
self.prioritized_messages.publish(context, message_to_prioritize.value)
例子: examples/python/dex_examples/primitives/channel/channel_flow.py
Transactional commit 不会在 handler 运行期间提供 isolation。如果决策依赖整个 loaded Channel、ChannelMap 或 AttributeMap 保持不变,所有协作的 Step 与 RPC writer 都必须使用同一把 Attribute lock。该 lock 提供 isolation,并隐式使 RPC 具备 transactional 语义;它不会加载 collection content。
Non-transactional RPC 遇到 missing deletion 时会把它当作 no-op,仍然 commit 其他成功的 effect。Cadence 的 Channel deletion validation 和 write 不具备相同的 atomic guarantee。Cadence caller 必须能容忍这项竞态。
RPC timeout
RPC timeout 限制一次 Worker handler invocation。
Python、Java、TypeScript 和 Rust 在 RPC handler 上声明 timeout。Go 在每次 Client 调用的 InvokeOptions 中传入。应用不设置 timeout 时,Dex Server 使用自己的 default。默认最大值为 60 秒。
@rpc(timeout=timedelta(seconds=30))
Attribute locks
Attribute lock 请参阅 Locking。