Skip to main content
Python Executor 是 Worker 本机 Function Executor Pool 里的一个常驻 Python 进程(python -m blockx_executor)。它从 Worker 收 ExecuteCall、装载并执行用户函数,把用户函数发起的 BlockDB / RPC / 子函数调用翻译成 CallWaiting 回送 Worker,收到 ResumeCall 后恢复原执行栈。它在系统中的位置见 架构总览,Worker 侧的对端(dispatcher、executor adapter、子函数)见 Call 执行子系统

职责与边界

负责:
  • 通过一条 UDS 长连接接收 ExecuteCall / ResumeCall / CancelCall / HeartbeatAck,发出 CallWaiting / CallCompleted / CallFailed / Heartbeat
  • 按源码 digest 编译并缓存函数模块,按 task 隔离模块命名空间,缺源码时按 digest 向 Worker 回拉。
  • greenlet 让同步风格的用户代码在 SDK 调用点挂起、在 ResumeCall 后恢复。
  • 在安全点执行取消与 deadline / budget 检查,并用 SIGALRM 定时器抢占纯 CPU 循环。
  • 每个心跳周期上报进程指标(RSS、CPU 利用率、累计 CPU 秒、各状态 call 计数)。
不负责:
  • 决定 call 派发顺序或选择哪个 Executor(Worker dispatcher 的事)。
  • 直接访问 BlockDB / 节点 RPC / Function Code View。SDK 的所有 IO 都被拦截后经 Worker 出去。
  • 持久化或跨进程恢复任何 call。进程重启后一切上下文重建。
核心不变量:
  • 用户代码只在主线程的 greenlet 里跑,每次只跑一个;UDS reader / writer / heartbeat 是独立线程,不切换 greenlet。
  • 一个 ExecContext 同一时刻最多一个 pending requestIdResumeCall 必须同时匹配 callId、状态对应的 resumeKindrequestId,否则按 stale 丢弃。
  • 同一 callId 再次 ExecuteCall 视为新 attempt:旧 greenlet 标 stale 并终止、不再发终态;新 attempt 用新 greenlet。
  • 每个被丢弃的 context 恰好发出一条终态消息(CallCompletedCallFailed),attempt 替换除外。
  • 编译产物按 digest 跨 task 共享;exec 出来的模块命名空间按 task_id 隔离,task 最后一个 call 结束时释放。

代码位置

核心类型与接口

Provider 协议原文:

用户函数编程模型

一个函数模块是一段 Python 源码,入口是模块顶层的一个 def。入口名解析规则在 blockx_audit/tables.pyresolve_entry_name:优先取 functionId 最后一段同名的顶层 def,其次是 def _(或 _ = some_def 别名),最后兜底第一个公开 def。实践中写 def _(...) 即可。 ExecuteCall.payload.args 是 JSON 数组时按位置展开传参,否则整体作为单参数(runtime.py _run_callabletarget(*input_payload) if isinstance(input_payload, list) else target(input_payload))。返回值直接是行数据:dict 表示单行,list[dict] 表示多行,None 表示无输出(详见 Plugin 系统)。返回值经 uds_codec.dumps_wire 序列化:超出 64 位的 Python int(如 uint256)可直接返回,datetime / date 转 RFC 3339 字符串,dataclass 经 asdict 兜底转换(并打一次告警),Decimal 等其他类型不可序列化,需先转成 str(下方示例里的 str(total) 就是这个原因)。 下面是 cmd/worker/worker_process_subcall_test.go 里的一个真实用例。它没有任何 import,因为非受限命名空间预置了 tables.PRELUDE_IMPORTSjsonDecimaldatetimedateblockdb.Table / BlockTable / TimeTable / Blockblockx.functionleafage.ChainState 等):
e2e/perf/s4_noderpc_batch_test.go 里的节点 RPC 读法:
各类 SDK 调用的用户侧写法与被拦截的位置: blockx_sdkblockx-py 的关系:blockx-py / blockdb-py / leafage-py 是用户真正面对的 SDK,也能在 Executor 之外直连服务;blockx_sdk 不是它们的替代品,而是 Executor 进程内的接缝——blockx_sdk.runtimecontextvars 保存当前 Providerblockx-pyfunction.call 在检测到 Executor 环境时 lazy import 它。blockx_sdk.db / rpc / call_subfunc 这组早期入口现在主要出现在单测、Go 侧 E2E fixture 和 cmd/bundle-loadgen 生成的负载函数里。
审计开启时(ExecuteCall.payload.audit 存在),模块在 restricted_runtime.restricted_namespace 里 exec:__builtins__ 只有 SAFE_BUILTINSimport 只能拿到白名单模块的 facade,blockx 只暴露 function 与能力类。哪些语法与 API 在白名单内见 blockx 仓库 docs/specs/function-code-python-whitelist.md

数据流 / 执行流程

主循环在 __main__.pywhile session.is_connected(): runtime.run_until_idle(); session.wait_for_message(timeout=runtime.next_scheduler_delay(None))run_until_idle 每轮先唤醒到期 sleeper、刷新连接状态、处理入站消息,再从 _sleep_ready(优先,连续 64 个后让位)和 _runnable 取一个 call_id 交给 _step_context_step_context 只有三个分支:有 cancel_reason_terminate_context;有 resume_payloadswitch(payload) 回等待点;否则首次 switch() 进入 _run_callable 启动时 ExecutorRuntime.__init__ 完成一次性 hook 安装:install_global_sleep_patch()_install_blockdb_bridgeblockdbmode.set_mock(False) + set_channel_factory)、_install_noderpc_bridge_install_localtest_bridgeblockx-py 的 channel factory,供能力后端的 Python 客户端用),并设置 BLOCKX_EXECUTOR_RUNTIME=1SdkBridgeProvider 不是全局装一次,而是每次进入 _run_callableset_provider、退出时 reset_provider
SDK hook 由两种手法组合而成:blockx_sdk 走 provider 注入;blockdb-pyblockx-py 通过它们暴露的 set_channel_factory 接缝在 gRPC channel 层被替换;leafage-py 通过替换 W3_DICT 条目被 patch。socket 层没有任何 hook。
BlockDB 与 NodeRPC 请求体不走 JSON:MessageEnvelope.binary_payload 挂 protobuf / JSON-RPC 原始字节,uds_codec.encode_frame 打成 BXB1 标签帧;Worker 回的 ResumeCall 也用 sidecar 带响应字节,runtime._handle_inbound_message_typed 把它塞进 result["Data"],bridge 再交给 gRPC stub 反序列化。blockdb-py 侧另有 mode / channel / stub 缓存。Executor 进程内没有 Python 本地 L1 读缓存,也没有协议批处理。

状态与生命周期

ExecState 的实际取值(context.py): 几点补充:
  • 恢复时状态不会先回到 RUNNABLEresume_call 只写 resume_payload 并置 queued=True 入队,_perform_waitswitch() 返回后直接改成 RUNNING。等待中的 context 是否已排队由 queued 标志表达,不由状态表达。
  • DONE 的 context 立即从 _contexts 移除,所以 inflight_count() 就是 len(_contexts)
  • WAITING_FUNCTION_CODE 用于按 digest 回拉源码。同一 digest 的并发 miss 在 SdkBridge._function_code_fetches 里做本地 singleflight,只有 leader 发 CallWaiting(waitKind=function_code)
attempt 与 requestId。 _replace_existing_attemptsubmit_call 里执行:旧 context 标 stale、从队列剔除、_terminate_context(..., suppress_terminal_emit=True)attempt_seq 来自 ExecuteCall.payload.attemptSeq(默认 1),不是本地递增。requestId 形如 req-<processInstanceId>-<seq>next_request_id),进程实例 id 前缀让重启后的进程不会与旧进程的 pending 请求撞号,seq 在进程内单调递增。stale resume / cancel 分别累加 _stale_resume_dropped / _stale_cancel_dropped 取消。 cancel_callWAITING_* 立即 _terminate_context(向 greenlet throw(CallCancelledError),再 GreenletExit);对 RUNNING / RUNNABLE 只设 cancel_reason 并入队,由下一个安全点的 check_cancelled 抛出。CancelCall.payload.reason 缺省为 worker_shutdown;本地超时闩用 call_timeout,终态归类为 deadline_exceeded 而不是 cancelled deadline 与 budget。 ExecContext.deadline_expired_reason() 是唯一判定点:budget_mscallBudgetMs,只计 CPU + IO 后端 + 子函数 + sleep,排队不计)优先于绝对的 call_deadline_ms。安全点有:进入用户函数前、每次 _perform_wait 发送前与恢复后、sleep 醒来后、返回前。此外 _switch_greenlet_with_deadline 在每次切入 greenlet 前 signal.setitimer(ITIMER_REAL, remaining, 0.1)_on_deadline_signal 只在该 call 处于 RUNNING 且确实过期时抛 _DeadlineSignalExpired。这条 SIGALRM 通路能抢占纯 CPU 循环;吞掉 BaseException 的循环或长 C 调用仍要靠 Worker 的 wedge 处理。 心跳与连接。 _heartbeat_loop 独立线程按 heartbeat_interval_sHeartbeat,payload 有 runnableCountinflightCountavailableContextsmax(0, max_inflight - inflight);不限时报 1<<30)、memoryBytesexecutorCpuUtilizationexecutorCpuSecondsTotalexecutorProcessInstanceIdlastSchedulerActiveAtMsrunningCallIdcallStateCountssentAtMsProcessMetricsSampler 在构造 ExecutorRuntime 时先做带重试的能力探测,失败则进程启动报错;运行期采样失败沿用上次值。UDS 断开后 _refresh_transport_state 给所有活动 context 打 transport_lost 并入队收敛,submit_call 拒绝新 call。 准入。 顶层 call 在 inflight_count() >= executor_max_inflight 时抛 ExecutorCapacityExhausted(回 CallFailed(errorKind=executor_capacity_exhausted, retryable=true));子函数不看该闸门,只要求父 context 驻留本进程,否则 subcall_parent_not_resident(不可重试)。executor_max_runnable 只是接收并保存,本地不强制。

配置

命令行参数(blockx_executor/__main__.py,由 internal/worker/adapters/executor/spawner.goexecutorRuntimeArgs 传入): 环境变量: 沙箱模式(EXECUTOR_SPAWN_MODE=sandbox)下 Go 侧只透传 internal/worker/adapters/executor/pool.gochildEnvKeys 白名单中的变量(上表各项加 LEAFAGE_ENDPOINTPROXY_SOCKET_PATH);裸进程模式仍继承全部父环境。容器空 netns,因此 OTLP 直连不可用。Python 进程本身对两种 spawn 模式无感,见 blockx 仓库 docs/specs/executor-sandbox-isolation.md

扩展点

  • 接一个新的外部 SDK 到 Worker IO 通道:写一个 bridge 模块(参考 blockdb_bridge.py / noderpc_bridge.py),把 SDK 的传输接缝替换成对 provider.read_req / write_req 的调用,在 ExecutorRuntime.__init__ 里安装;backend 名要与 Worker 侧 IO 子系统注册的 backend 对齐(见 IO 访问子系统)。若 SDK 是 gRPC unary,直接复用 BridgeChannel(provider_getter, backend=..., streaming_methods=frozenset())
  • 给用户代码新增 SDK 入口:只加 blockx_sdk 下的薄封装并让它调 get_provider();不要在 SDK 里持有连接。受限模式还要同步 blockx_audit/tables.py 的白名单与 restricted_runtime facade。
  • 新增 ResumeCall 错误类型:在 sdk_bridge.pyRemoteCallError 子类并扩展 _map_resume_error;backend 专属映射用 register_io_error_kinds
  • 新增心跳字段或消息类型messages.py 加构造、wire_keys.py 加常量、runtime._handle_inbound_message_typed 加分支,并同步 api/uds/types.go
  • 改安全点或计费:所有过期判断经 ExecContext.deadline_expired_reasonraise_if_deadline_expired,不要在别处另写判断。

测试

python/tests/ 按实现边界分文件(test_<boundary>.py),文件内按逻辑块分 TestCase<Boundary><Theme>Test),规范见 blockx 仓库 docs/specs/python-test-organization.md。主要边界:
  • Runtime:test_executor_runtime_scheduler.pytest_executor_runtime_termination.pytest_executor_runtime_attempt.pytest_executor_budget.pytest_executor_subcall.pytest_executor_oom_hardening.pytest_terminal_bounds.py
  • 协议:test_executor_protocol_contract.pyExecuteCall 校验、入站控制消息、心跳字段)、test_uds_codec.pytest_uds_session.pytest_envelope_*.py
  • SDK hook:test_executor_sdk_hook.pytest_sdk_bridge_provider.pytest_sdk_bridge_capability.pytest_perform_wait_emit_failure.pytest_blockdb_bridge.pytest_noderpc_bridge.pytest_sleep_patch.py
  • 装载与隔离:test_module_registry.pytest_restricted_runtime.pytest_executor_audit_gate.py
  • 上下文与指标:test_exec_context_*.pytest_process_metrics_sampler.pytest_logging_setup.pytest_tracing_setup.pytest_debug_capture.py
  • 审计器:test_audit_*.py
单测里不需要真实 UDS:ExecutorRuntime("exec-1") 不传 session 时消息进 drain_outbox(),用 handle_inbound_message(dict) 喂入站消息,run_until_idle() 推进。 E2E 在 Go 侧:cmd/worker/worker_process_*_test.gocmd/syncinvoker/*_test.go 拉起真实 Worker + python/.venv/bin/python -m blockx_executore2e/system 起整套系统;e2e/perf 用仓库根 test.shsystemd-run cgroup 限额下跑。命令见 测试组织与命令

相关文档

blockx 仓库 spec:
  • docs/specs/python-function-executor.md — Executor 进程内设计:调度、状态、attempt / requestId、取消与 deadline、资源边界(重点)。
  • docs/specs/worker-executor-connection-and-python-sdk-hook.md §6–§8 — SDK hook 方案、SdkBridge 同步 contract、挂起恢复时序。
  • docs/specs/python-executor-process-metrics.md — 心跳里进程指标的采样与降级。
  • docs/specs/2026-07-23-execute-call-function-code-reference.md — 按 digest 回拉源码的协议。
  • docs/specs/executor-sandbox-isolation.md — containerd 沙箱 spawn 模式。
  • docs/specs/2026-07-30-python-executor-cpu-optimization.md — CPU 优化方案。
  • docs/specs/python-test-organization.md — 单测组织标准。
  • docs/specs/function-code-python-whitelist.mddocs/specs/function-code-audit-design.md — 用户代码白名单与审计。
站内: