跳到主要内容

Channel

Channel 是一条单个 Flow execution 内部 durable 的 first-in, first-out (FIFO) 队列。Step 与 RPC 能往里面追加有类型的消息。Step 通过在 WaitFor 里返回 Channel condition 来消费消息。

Channel 会按到达顺序保留消息。一条消息只能被一个匹配的 wait 消费一次,不会复制给其他 waiting Step。

Channel 只属于一个 Flow execution。消息不会跨 Flow execution。

Persistence schema

Channel 有稳定的名字和消息类型。先定义它,再加到 Flow 的 persistence schema。

下面的例子有一个 ApprovalMessages Channel。start Step 等一条 approval 或一个 Timer。publishApprovalMessage RPC 发布 approval 消息。应用通过 enqueueChannelMessage RPC 发布另一条 queued message。

@dataclass(frozen=True)
class QueuedMessageReference:
message_id: str


class ChannelWaitStep(Step[int]):
def __init__(
self,
approval_messages: Channel[str],
queued_messages: Channel[str],
) -> None:
self.approval_messages = approval_messages
self.queued_messages = queued_messages

def get_step_options(self) -> StepOptions:
return StepOptions(execute_load_channels=(self.queued_messages,))

def wait_for(self, context: Context, input: int) -> Wait:
return Wait.any_of(
self.approval_messages.for_one(),
Timer.by_duration(timedelta(seconds=input)),
)

def execute(self, context: Context, input: int) -> StepDecision:
pending_queued_messages = self.queued_messages.pending_messages(context)
if pending_queued_messages:
self.queued_messages.delete(context, pending_queued_messages[0].message_id)
return graceful_complete(pending_queued_messages[0].value)
if context.has_timer_fired():
return graceful_complete("approval timed out")
approval_message_values = self.approval_messages.results(context)
return graceful_complete(approval_message_values[0])


class ChannelFlow(Flow[int]):
approval_messages = Channel("ApprovalMessages", str)
queued_messages = Channel("QueuedMessages", str)
prioritized_messages = Channel("PrioritizedMessages", str)

def __init__(self) -> None:
self.wait_for_approval = ChannelWaitStep(self.approval_messages, self.queued_messages)

def get_steps(self) -> StepList[int]:
return StepList.start_step(self.wait_for_approval)

def get_persistence_schema(self) -> PersistenceSchema:
return PersistenceSchema.of(self.approval_messages, self.queued_messages, self.prioritized_messages)

@rpc
def publish_approval_message(self, context: Context) -> None:
self.approval_messages.publish(context, "approved")

@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

Timer 分支会先返回,不会读取 Channel 结果。approval 胜出时,Channel condition 已消费一条消息,所以第一个结果一定存在。

等待并消费

WaitFor 里返回 Channel condition 会等待这个 Channel。condition 满足后,它会从 Channel 消费消息。Dex Server 会保留消息,直到某个 waiting Step 赢得它们;然后按 FIFO 顺序从队列移除选中的消息,并以 Channel results 传给 Execute

Condition何时继续消费什么
ForOne至少有一条消息。恰好一条消息。
ForN至少有指定数量的消息。恰好指定数量的消息。
AtLeast至少有指定数量的消息。当时所有可用的消息。
AtMost立即继续,包括队列为空时。当前已排队消息中的最多指定数量。
AtLeastAtMost至少有下界指定数量的消息。当前已排队消息中从下界到上界的数量。

AtMost 本身不会等待消息累积。外层 Wait 完成时,它会消费当时队列中不超过上限的消息。要等待至少一条消息并消费不超过指定上限的消息,请使用下界为 1 的 AtLeastAtMost

发布消息

在 Step 或 RPC 内,用 handler 的 Context 调 Channel 的 publish method。Dex 会在接受这个 handler result 时追加消息。

Worker 外的代码调用 Flow RPC,由 handler 使用 Context 发布。这样 Flow 就能统一处理 validation、authorization、lock 和相关状态变更。

消息可以先于 Step 的 wait 到达。Dex Server 会把它保留在 Channel 里,直到之后的 wait 消费它。

管理 pending message

Dex 接受每次 publication 时都会分配一个 UUIDv7 message ID。RPC 可以按 FIFO 顺序列出所有 pending message,也可以用 ID 删除其中一条。List 不会消费 message。只有仍然 pending 的 message 才能删除;如果并发 Step 已经消费它,或者另一个请求已经删除它,API 会返回 Channel-message-not-found error。

这是 queue state,不是 Flow history。Step 消费 message 后,它会从 Channel 消失,但 publication 与 consumption 仍可在 Flow history 中观察,直到 retention 将它们清理。

每个 Step method、timeout handler 和 RPC 始终会收到 Channel size metadata。读取 pending message 的 ID 或 Value 前,必须显式加载对应 Channel。删除 ID 已知的 message 不需要加载它。StepOptionsWaitForExecute 分别选择状态;FlowTimeoutHandlerOptions 为 timeout handler 选择状态。加载会建立一份 invocation snapshot,但不会消费 message。已加载的空 Channel 返回空列表;未加载就读取 pending message 会返回 usage error。

应把 Step 或 RPC 中读取的每一份 pending-message snapshot 都视为可能已经过期。handler 运行时,其他 Step 和 RPC 仍可能消费、删除或发布消息。transactional execution 会校验选中的 deletion,并原子提交本次 handler 的 write,但不会锁定整个 Channel snapshot。只有当该操作明确能够容忍这种并发时,才应直接读写 pending message。如果业务判断要求队列保持不变,所有会写入它的 Step 和 RPC 都必须协作使用同一个 Attribute lock。

Step 或 timeout handler 可以把 deletion 与其他 side effect 一起暂存。成功的 handler effect 按 Attribute write、best-effort Channel deletion、Channel publication 的顺序应用,然后处理 Step decision。如果另一项操作已经消费或删除所选消息,deletion 是 no-op,其他 effect 仍会应用。StepDecision 只描述 control flow,不包含 deletion。

RPC 可以把删除与其他 durable write 组合。Move 示例加载 QueuedMessages,按 message ID 找到原始 Value,再暂存 deletion 并重新发布该 Value。Caller 只传 ID。Missing message 必须中止所有 write 时,需要启用 transactional execution;见 Transactional read 与 write

@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

ChannelMap

一个 Flow 有一组数量不受限、彼此独立的队列时,用 ChannelMap,例如每个订单、客户或 job 一条 Channel。ChannelMap 只声明一个名字和一个值类型,每个 instance key 都有自己的 FIFO 队列。

RPC 无需加载 pending message,也能读取当前 ChannelMap key 和 size。Handler 需要所有 instance 的 pending message 时,加载整个 ChannelMap;只需要已知 key 时,加载指定 instance。Instance key 不能包含 /,因为 Dex 使用该字符分隔 map name 与 instance key。