任务执行:Worker 与 RunContext
30 秒导读: 第 2 章的 Executor 是"大脑",它决定某个任务该跑了,就往队列里丢一个
WorkerJob。本章讲的 Worker 是"手脚":它把这个 job 从队列里拉过来,在一个隔离的执行环境(RunContext)里真正调用插件代码RunnableTask.run(),再把结果发回去。Worker 有一条铁律——它从不直接碰数据库,一切外部状态都走 gRPC 门面(MetaStore/StateStore)。
1. 这是什么(零基础也能懂)
一句话定义: Worker 是 Kestra 里真正执行任务代码的进程。
Executor 负责"推进状态机"(第 2 章),但它自己不跑任务——跑一个 Python 脚本、发一个 HTTP 请求、查一次数据库,这些活儿都在 Worker 上发生。这样设计是为了把"编排"和"执行"拆开:Executor 轻量、快、只碰状态;Worker 重、可能长时间阻塞、可能崩,但崩了不影响编排大脑。
给谁用 / 解决什么问题: 想象你写了一个 flow,里面有个任务 runScript。当轮到它跑时:
- Executor 把它包装成一个
WorkerTask(一种WorkerJob)扔进队列。 - 某台 Worker 把它领走,在自己机器上建一个临时工作目录、渲染好所有
{{ }}表达式、调用插件的run()。 - 插件返回 outputs,Worker 把
WorkerTaskResult发回给 Executor,状态机继续推进。
核心直觉/类比: 把 Executor 当餐厅前台(记单、叫号、更新订单状态),Worker 是后厨的厨师(真正炒菜)。前台永远不进厨房,厨师也永远不改订单系统——两者之间只通过"传菜口"(gRPC 流)递条子。
本章覆盖三大块:
| 块 | 讲什么 | 关键类 |
|---|---|---|
| 分发骨架 | job 怎么从 controller 流到 worker | WorkerJobDispatcher、WorkerJobFetcher |
| 作业循环 | worker 内部怎么逐个跑 job | WorkerLoop、WorkerJobExecutor、WorkerTaskProcessor |
| 执行沙箱 | 插件在什么环境里跑 | RunContext、DefaultRunContext |
本章不讲状态机推进(见第 2 章)和触发器如何产生 job(见第 4 章);队列抽象本身见第 5 章。
2. 顶层全景(它大概怎么转)
先看一次任务从"被派发"到"结果回传"的完整数据流。怎么读这张图:从左到右是一次 job 的生命周期,虚线是控制信号(permits/completion),实线是数据(job/result)。
controller 侧 (worker-controller) worker 侧 (worker)
┌───────────────────────────────────┐ ┌──────────────────────────────────┐
│ WorkerJobDispatcher │ │ WorkerJobFetcher (一条 gRPC 双向流) │
│ · 订阅 job 队列 │◀──permits/completion(虚线)──│ │
│ · 按 permits+bucket 选一个 worker │ │ · 拉 job 放进本地 WorkerQueue │
│ · 持久化到 StateStore 再发 │──job──▶│ │
└───────────────────────────────────┘ └──────────────┬───────────────────┘
▲ │ poll
│ WorkerTaskResult ▼
┌─────────────────┴─────────────────┐ ┌──────────────────────────────────┐
│ GrpcWorkerControllerService │◀─send──│ WorkerJobExecutor │
│ · sendWorkerTaskResults→结果队列 │ │ · N 个 Consumer(虚拟线程)poll │
│ · sendWorkerLogEntries/Metrics │ │ · 每个 job 提交到平台线程池执行 │
└────────────────────────────────────┘ │ │ │
│ ▼ │
│ WorkerTaskProcessor │
│ → WorkerTaskCallable │
│ → task.run(RunContext) ★ │
└──────────────────────────────────┘
★ = 真正的插件代码在这里跑
部件一句话职责:
| 部件 | 干什么 | 文件 |
|---|---|---|
WorkerJobDispatcher | controller 侧总协调:选 worker、占 bucket、发 job | worker-controller/.../services/WorkerJobDispatcher.java |
GrpcWorkerControllerService | gRPC 服务端:开双向流、收 permits、收各类结果 | worker-controller/.../services/GrpcWorkerControllerService.java |
WorkerJobFetcher | worker 侧拉 job 的那条流的客户端 | worker/.../fetchers/WorkerJobFetcher.java |
WorkerJobExecutor | worker 侧线程池 + N 个消费者,逐个跑 job | worker/.../WorkerJobExecutor.java |
WorkerTaskProcessor | 把一个 WorkerTask 变成真正的执行:建 RunContext、重试、缓存 | worker/.../processors/WorkerTaskProcessor.java |
WorkerTaskCallable | 最内层:真正调用 RunnableTask.run(),处理超时 | worker/.../processors/internals/WorkerTaskCallable.java |
RunContext | 执行沙箱:变量渲染、存储、工作目录、日志、密钥 | core/.../runners/RunContext.java |
GrpcWorkerIOSender | worker 侧把 result/log/metric 发回 controller | worker/.../senders/GrpcWorkerIOSender.java |
主线走一遍(高层):
- Worker 启动时,
WorkerJobFetcher开一条到 controller 的 gRPC 双向流,报上自己的 id、group、最大并发数,以及初始 permits(可接活余量)。 WorkerJobDispatcher一旦从队列拿到 job,就找一个有 permit 的 worker,先把 job 持久化到 StateStore(为了崩溃恢复),再通过流发过去。- Worker 收到后放进本地
WorkerQueue;WorkerJobExecutor的某个消费者 poll 到它,提交给线程池。 WorkerTaskProcessor建好RunContext,WorkerTaskCallable调用task.run(runContext)——插件代码在这一刻执行。- 结果、日志、指标各走一条
GrpcWorkerIOSender发回 controller;job 完成时 worker 还会回传一个 completion 信号,让 controller 释放占用的容量。
3. 分发骨架:job 怎么从 controller 流到 worker
这节讲任务与执行器之间的 gRPC 分发。核心难点是:controller 有很多 job,worker 有有限容量,怎么做到既不压垮 worker、又不丢 job。Kestra 的答案是基于 permits 的流控 + pull/ack 模式。
3.1 一条长连接,双向流
Worker 侧只开一条 gRPC 双向流(streamWorkerJobs),它同时承载三种信息:
| 方向 | 内容 |
|---|---|
| worker → controller | permits(剩余容量)、completion(哪些 job 到终态了) |
| controller → worker | job payload、广播事件(kill 命令、集群事件、元数据变更) |
流的客户端是 WorkerJobFetcher——注意它本身是一个 WorkerLoop(见 §4.1),循环体 doOnLoop() 负责建流、按需重连(指数退避)、周期性上报 permits(worker/.../fetchers/WorkerJobFetcher.java:258 doOnLoop、:293 startStream)。服务端在 GrpcWorkerControllerService.streamWorkerJobs 里为每条流建一个 WorkerStreamContext,并把它注册进 dispatcher(worker-controller/.../services/GrpcWorkerControllerService.java:86)。
3.2 permits:worker 说"我还能接几个"
permits 就是 worker 本地队列的剩余容量——一个数字,代表"你还可以再发几个 job 给我"。计算很直白:
// 示意,非源码:calculatePermits 的核心
private int calculatePermits() {
return fetchingPaused.get() ? 0 : workerJobQueue.remainingCapacity(); // 维护模式下报 0
}
真实实现在 WorkerJobFetcher.calculatePermits(worker/.../fetchers/WorkerJobFetcher.java:519)。关键点:permits 是一个电平值(level)而非增量,所以 0 是有意义的——它表示"别再发了"。维护模式/cordon 时 worker 就上报 0,让 controller 停止派发,但流保持存活(还能收 kill 命令)。
3.3 dispatcher 侧:选谁、占坑、发货
controller 从队列拿到一个 job 后,handleIncomingJob 会做三件事(worker-controller/.../services/WorkerJobDispatcher.java:711):
- 预检 kill: 如果这个 execution 已经被 kill,直接产出一个
KILLED结果,不派发。 - 找 worker 并占坑:
findAndReserveWorker挑一个"最闲"(in-flight 最少)的候选,先 CAS 消费一个 permit,再占一个 bucket 槽位(:852)。两步都成功才算占到;bucket 占不到就把 permit 还回去,试下一个候选。 - 发货:
dispatchJobToWorker先持久化到WorkerJobRunningStateStore再 sendResponse(:888)。这个顺序是崩溃恢复的关键——只要 job 要么在队列里、要么在 StateStore 里,controller 崩了也不会丢。
permit 和 bucket 是两个独立的闸门,为什么要两个?
- permit = worker 的总余量(所有队列共享),防止过量派发。
- bucket = 某个 Worker Queue 的容量槽,由
WorkerCapacityPolicy管理,支持按队列预留容量(worker-controller/.../services/WorkerCapacityPolicy.java:27)。
默认策略 SinglePoolCapacityPolicy 最简单:一个 maxConcurrency 大小的共享池,占坑就是原子自增(worker-controller/.../services/SinglePoolCapacityPolicy.java:30 tryReserve)。需要"给某队列保底 N 个槽"的部署,替换 WorkerCapacityPolicyFactory 这个 bean 即可。
3.4 completion:占的坑什么时候还
一个微妙点:bucket 槽位从派发一直占到 job 真正跑完,而不是"worker 一收到就还"。job 到终态时,worker 侧的消费者会调 workerJobFetcher.onJobCompleted(jobId)(worker/.../WorkerJobExecutor.java:322),把 jobId 塞进 completion 队列、并立刻 flush 一次(worker/.../fetchers/WorkerJobFetcher.java:556)。controller 收到后在 onCompletionsReceived 里释放对应 bucket,并重新评估该队列的 pause/resume(WorkerJobDispatcher.java:668)。
为什么不等下一个心跳周期? completion 立刻 flush 能省掉最多一个
PERMIT_CHECK_INTERVAL(100ms)的空窗,让"同一 bucket 门控的下一个任务"更快派发出去。
4. Worker 作业循环:job 到了之后怎么跑
这节讲 worker 进程内部:从"本地队列里有个 job"到"插件代码被调用"这一段。这是层层嵌套的四层结构,每层职责单一。
4.1 底座:WorkerLoop —— 一个可暂停/可优雅停的循环
Worker 里所有"后台循环"(拉 job 的 fetcher、发结果的 sender、跑 job 的 consumer)都继承 WorkerLoop。它把"循环"这件事的通用骨架抽出来:
// 示意,非源码:WorkerLoop 的骨架
while (running.get()) {
waitIfPaused(); // 支持暂停(维护模式)
if (!running.get()) continue;
doOnLoop(); // 子类实现:干一次活
// 抛 InterruptedException → 退出;抛别的 → 记日志继续
}
cleanup(); // finally:释放资源
真实实现在 worker/.../WorkerLoop.java:88(runLoop)。子类只需实现 doOnLoop()(:121),外加可选的 signalJobStop()(:238,优雅停机时通知正在跑的活)。一个抽象类,三种用途(fetch / send / consume),是 worker 模块最重要的复用点。
4.2 分发线程:WorkerJobExecutor 与虚拟线程消费者
WorkerJobExecutor.start() 会起 N 个 WorkerJobConsumer,跑在虚拟线程上(worker/.../WorkerJobExecutor.java:82)。为什么是虚 拟线程?因为消费者本身只做"poll + 等待",几乎不吃 CPU,虚拟线程廉价。
但真正跑任务的却是平台线程池——消费者 poll 到 job 后,把它 submit 给一个 maxCachedThreadPool,然后自己 future.get() 阻塞等它跑完:
// 示意,非源码:WorkerJobConsumer.doOnLoop 的核心
WorkerJob job = workerJobQueue.poll(Duration.ofSeconds(1));
if (job == null) return;
WorkerJobProcessor processor = factory.create(context, job);
Future<?> future = taskExecutorService.submit(() -> processor.process(job)); // 提交到平台线程池
future.get(); // 消费者在此阻塞,直到任务跑完
workerJobFetcher.onJobCompleted(job.uid()); // finally:回传 completion
真实实现在 WorkerJobExecutor.WorkerJobConsumer.doOnLoop(worker/.../WorkerJobExecutor.java:283)。
为什么消费者(虚拟线程)和任务(平台线程)要分开? 这是中断隔离:任务超时/被 kill 时,只中断跑任务的那个平台线程,消费者的虚拟线程不受影响,能干净地继续 poll 下一个。任务里的异常也被
future.get()完全兜住,消费者永远不会因为任务出错而挂掉(:310)。
4.3 处理器:WorkerJobProcessor 家族
WorkerJobProcessor 是一个接口,只有三个方法:process() / stop() / kill()(worker/.../processors/WorkerJobProcessor.java)。工厂按 job 类型选具体实现:
| job 类型 | 处理器 | 干什么 |
|---|---|---|
WorkerTask(普通任务) | WorkerTaskProcessor | 本章重点 |
WorkerTrigger(轮询触发) | WorkerTriggerProcessor | 见第 4 章 |
AbstractWorkerJobProcessor 用一对 CAS 保证一个处理器同时只跑一个 job,并管理 kill/stop 信号的转发(worker/.../processors/AbstractWorkerJobProcessor.java:52)。有个精巧的竞态处理:如果 kill 信号在 callable 被记录之前就到了,callJob 会在设置 callable 后补一次 kill(:69),防止 kill 丢失。
WorkerTaskProcessor.doProcess 分三种情况(worker/.../processors/WorkerTaskProcessor.java:98):
RunnableTask→runTask(),正常执行一个任务。WorkingDirectory(流程控制任务)→runWorkingDirectory(),在 同一个工作目录里顺序跑一串子任务,子任务间用runIf决定跳过。- 其它 → 报错(worker 只跑 runnable task)。
runTask 里还夹着输出缓存逻辑:如果任务开了 taskCache,先按渲染后的任务定义算 SHA-256 哈希查缓存,命中就跳过执行直接返回缓存的 outputs(:225)。
4.4 最内层:WorkerTaskCallable —— 真正调 run() 的地方
runAttempt 会先发一个"任务进入 RUNNING"的结果(让 Executor 知道任务动了),然后构造 WorkerTaskCallable 并执行(worker/.../processors/WorkerTaskProcessor.java:362、:396)。
WorkerTaskCallable.doCall 就是那条"最深的线":
// 示意,非源码:WorkerTaskCallable.doCall 的核心
Duration timeout = runContext.render(workerTask.getTask().getTimeout())...;
if (timeout != null) {
Failsafe.with(Timeout.builder(timeout).withInterrupt().build())
.run(() -> taskOutput = task.run(runContext)); // ★ 插件代码在这里跑
} else {
taskOutput = task.run(runContext); // ★
}
return taskOutput.finalState().orElse(SUCCESS);
真实实现在 worker/.../processors/internals/WorkerTaskCallable.java:60。RunnableTask.run(RunContext) 就是所有插件的执行入口(core/.../models/tasks/RunnableTask.java:14)——整章的"手脚"最终落在这一行 task.run(runContext) 上(:84)。
它的父类 AbstractWorkerCallable.call() 有个关键决定:它 catch Throwable(不只是 Exception),因为插件可能抛 Error(依赖问题、坏行为),必须保证任务无论如何都会以某个状态收场,而不是让 worker 线程裸崩(worker/.../processors/internals/AbstractWorkerCallable.java:63)。超时/kill 通过中断当前线程实现(:109 kill)。
5. RunContext:执行沙箱
task.run(runContext) 里那个 runContext 是插件与 Kestra 世界的唯一接口。插件不能随便碰文件系统 、数据库、密钥——它能做的一切都通过 RunContext 暴露。这节讲这个沙箱提供了哪些能力。
RunContext 是抽象类(core/.../runners/RunContext.java:24),实际用的是 DefaultRunContext(core/.../runners/DefaultRunContext.java:53)。它在 worker 侧由 RunContextInitializer.forWorker 构造:从 wire 数据重建变量、装上存储、日志、工作目录、插件配置(core/.../runners/RunContextInitializer.java:85、:122)。
四大能力如下。
5.1 变量与 Pebble 渲染
flow 里到处是 {{ execution.id }}、{{ inputs.name }} 这种模板。渲染就是把这些占位符替换成真实值。
变量从哪来? RunVariables.DefaultBuilder.build 把一次执行的所有上下文组装成一个大 map:flow、execution、taskrun、inputs、outputs、vars、labels、envs… (core/.../runners/RunVariables.java:359)。所有结构化可用路径显式列在 EXECUTION_CONTEXT_PATHS(:49),还有测试专门防止它和实际 build 漂移。
谁来渲染? VariableRenderer 包着 Pebble 模板引擎(core/.../runners/VariableRenderer.java:22)。DefaultRunContext.render(inline) 只是转发给它(DefaultRunContext.java:259)。有个小优化:如果字符串里根本没有 {,直接短路返回,不进 Pebble(VariableRenderer.java:77)。默认还支持递归渲染(渲染 结果里若还有 {{ }} 继续渲染),上限 MAX_RENDERING_AMOUNT = 100 防死循环。
5.2 存储与工作目录
每个任务跑在一个独立的临时工作目录里,跑完清掉。这由 LocalWorkingDir 实现(core/.../runners/LocalWorkingDir.java:32)。
沙箱的边界就在 resolve():它保证任何要访问的相对路径解析后仍在工作目录内,否则抛异常——防止插件用 ../../ 逃逸出去(LocalWorkingDir.java:86):
// 示意,非源码:LocalWorkingDir.resolve 的安全检查
Path resolved = baseDir.resolve(path).toAbsolutePath();
if (!resolved.startsWith(baseDir)) {
throw new IllegalArgumentException("必须是工作目录内的相对路径");
}
FilesService 在此之上提供 inputFiles(把 {{ }} 渲染后的内容写进工作目录) 和 outputFiles(按 glob 收集产物上传到内部存储)两个静态门面(core/.../runners/FilesService.java:22、:74)。内部存储(跨任务传大文件)通过 runContext.storage() 暴露,worker 侧装的是 InternalStorage(RunContextInitializer.java:156)。
5.3 日志
runContext.logger() 返回的不是普通 logger,而是 RunContextLogger 装配的一套 logback appender 链(core/.../runners/RunContextLogger.java:34)。每条日志会被复制到三处:
| appender | 去向 |
|---|---|
ContextAppender | 打包成 LogEntry 发到日志队列(UI 上能看的那些行) |
FileAppender | 若任务配了 logToFile,写进可下载的日志文件 |
ForwardAppender | 转发到服务器本地日志 |
5.4 密钥(secret 遮蔽)
这是安全上最要紧的一环:密钥绝不能出现在日志里。机制分两步:
- 登记: 用到的密钥通过
RunContextLogger.usedSecret登记(明文 + Base64 两种形式都存),DefaultRunContext.setLogger会自动把 SECRET 类型的 input 登记进去(RunContextLogger.java:178、DefaultRunContext.java:182)。 - 遮蔽: 每条日志经过 appender 时,
BaseAppender.replaceSecret递归扫描消息和参数,把登记过的密钥替换成******(RunContextLogger.java:347)。
加解密本身由 Secret 这个包私有小类处理(core/.../runners/Secret.java:12);没配加密 key 时它会告警并返回原文而非崩溃(:22)。runContext.decrypt/encrypt 只是转发给它(DefaultRunContext.java:371)。
任务跑完,runContext.cleanup() 清掉工作目录并重置日志 MDC——此后不应再用这个 logger(DefaultRunContext.java:530)。
6. 结果回传:senders
任务产出的东西不止一种,Kestra 为每种开一条独立的 sender(都是 WorkerLoop 子类,worker/.../senders/GrpcWorkerIOSender.java:36),背后策略不同:
| sender | 内容 | 策略 | 失败重发? |
|---|---|---|---|
| taskResultSender | WorkerTaskResult | 逐条(PER_ITEM) | 是 |
| logEntrySender | LogEntry | 批量(BATCH) | 否 |
| metricsSender | MetricEntry | 批量(BATCH) | 否 |
差异的道理:任务结果丢了会让 execution 永远卡在 RUNNING,必须重发;日志/指标是高频、尽力而为的,重发反而可能给 worker 造成背压(GrpcWorkerIOSenderFactory.java 各 *Sender 方法)。onSendFailed 只对可重试的传输错误(UNAVAILABLE/DEADLINE_EXCEEDED)且非停机中才把批次重新入队(GrpcWorkerIOSender.java:201)。
结果流回 controller 侧的 GrpcWorkerControllerService.sendWorkerTaskResults,它 emit 到结果队列,并在任务到终态时删掉 StateStore 里的 running 记录(worker-controller/.../services/GrpcWorkerControllerService.java:194)。
7. 一条铁律:Worker 不依赖 repository
Worker 从不直接访问数据库。这是刻意的架构约束——Worker 可能是无状态、动态扩缩、甚至部署在隔离网络里的进程,给它数据库连接既不安全也不现实。
那 Worker 需要读 flow 定义、KV、命名空间文件怎么办?走 gRPC 门面。核心侧定义了 FlowMetaStoreInterface 这类接口,worker 侧用 gRPC 实现替换默认实现:
@Replaces(DefaultFlowMetaStore.class) // 用 gRPC 版顶掉直连 DB 的默认版
public class GrpcWorkerFlowMetaStore implements FlowMetaStoreInterface, ... {
private final WorkerFlowMetaStoreServiceBlockingStub workerFlowMetaStoreStub; // 只有 gRPC stub
}
真实实现在 worker/.../stores/GrpcWorkerFlowMetaStore.java:50,同类还有 GrpcWorkerKVMetadataStateStore、GrpcWorkerNamespaceFileMetadataStateStore。它们全都只持有 gRPC stub,没有任何 repository 依赖——需要数据时向 controller 发一次 gRPC 请求,由 controller 代查。这样 worker 对存储层完全无知,RunContext 里 namespaceKv() 等能力也顺着这条门面工作(RunContext.java:187)。
8. 巧妙之处(可借鉴的技术)
-
permit + bucket 双闸门,且占坑跨越整个 job 生命周期。 permit 管总量、bucket 管分队列配额;bucket 从派发一直占到 completion 回传才释放,保证"派出去但还没跑完"的 job 也算进容量。占坑失败时先还 permit 再试下一个候选,避免误暂停整个队列订阅(
WorkerJobDispatcher.java:852)。 -
虚拟线程 poll + 平台线程执行的中断隔离。 消费者廉价(虚拟线程),任务可被单独中断(平台线程),任务异常被
future.get()完全兜住,消费者永不因任务出错而死(WorkerJobExecutor.java:299)。 -
先持久化再派发。 dispatcher 在 sendResponse 之前必落 StateStore,保证 job 任何时刻要么在队列、要么在 StateStore,controller 崩了也不丢(
WorkerJobDispatcher.java:888)。 -
catch Throwable 兜底插件。 插件可能抛 Error,worker 宁可把任务标记为 FAILED 也 绝不让执行线程裸崩(
AbstractWorkerCallable.java:74)。 -
secret 登记 + 递归遮蔽。 密钥两种编码都登记,日志经 appender 时递归替换成
******,连嵌套 map/list 里的都不放过(RunContextLogger.java:347)。
9. 边界与局限
-
Worker 不推进状态机。 它只产出
WorkerTaskResult,状态怎么流转是 Executor 的事(第 2 章)。 -
只跑 runnable task。 非
RunnableTask/WorkingDirectory的任务 worker 直接报错;真正的"流程控制"(分支、循环)在 Executor 侧展开。 -
结果丢失的兜底不是无限的。 taskResult 会重发,但过了 liveness 超时 worker 会自断连接,由 executor 重新派发它的任务;重发只在运行中的循环内有界进行,停机时不重发以免
cleanup()排空变成死循环(GrpcWorkerIOSender.java:201)。 -
超大 payload 会被直接判失败。 如果一个 job 序列化后超过 worker 声明的 gRPC 入站上限,dispatcher 不发它、而是干净地把任务判 FAILED,避免 worker 反复 RESOURCE_EXHAUSTED、执行永久卡住(
WorkerJobDispatcher.java:888内rejectOversizedJob)。
10. 代码地图(导航索引)
| 主题 | 文件 | 关键符号 |
|---|---|---|
| 循环骨架 | worker/.../WorkerLoop.java | WorkerLoop、runLoop、doOnLoop、signalJobStop |
| 默认 worker | worker/.../WorkerAgent.java | WorkerAgent、resolveWorkerGroupId |
| 生命周期装配 | worker/.../AbstractWorker.java | AbstractWorker、start |
| 线程池 + 消费者 | worker/.../WorkerJobExecutor.java | WorkerJobExecutor、WorkerJobConsumer.doOnLoop |
| 处理器接口 | worker/.../processors/WorkerJobProcessor.java | WorkerJobProcessor、process、stop、kill |
| 处理器基类 | worker/.../processors/AbstractWorkerJobProcessor.java | AbstractWorkerJobProcessor、callJob |
| 任务处理器 | worker/.../processors/WorkerTaskProcessor.java | WorkerTaskProcessor、doProcess、runTask、runAttempt |
| 最内层执行 | worker/.../processors/internals/WorkerTaskCallable.java | WorkerTaskCallable、doCall |
| callable 基类 | worker/.../processors/internals/AbstractWorkerCallable.java | AbstractWorkerCallable、call、kill |
| 拉 job 的流 | worker/.../fetchers/WorkerJobFetcher.java | WorkerJobFetcher、calculatePermits、onJobCompleted |
| 结果 回传 | worker/.../senders/GrpcWorkerIOSender.java | GrpcWorkerIOSender、send、onSendFailed |
| 分发协调 | worker-controller/.../services/WorkerJobDispatcher.java | WorkerJobDispatcher、handleIncomingJob、findAndReserveWorker、dispatchJobToWorker |
| gRPC 服务端 | worker-controller/.../services/GrpcWorkerControllerService.java | GrpcWorkerControllerService、streamWorkerJobs、sendWorkerTaskResults |
| 容量策略 | worker-controller/.../services/WorkerCapacityPolicy.java、SinglePoolCapacityPolicy.java | WorkerCapacityPolicy、tryReserve、hasCapacity |
| 流上下文 | worker-controller/.../services/WorkerStreamContext.java | WorkerStreamContext、tryReserveBucket、completeJob |
| 沙箱抽象 | core/.../runners/RunContext.java | RunContext、render、storage、workingDir、logger |
| 沙箱实现 | core/.../runners/DefaultRunContext.java | DefaultRunContext、render、cleanup |
| 沙箱构造 | core/.../runners/RunContextInitializer.java | RunContextInitializer、forWorker |
| 变量组装 | core/.../runners/RunVariables.java | RunVariables、DefaultBuilder.build、EXECUTION_CONTEXT_PATHS |
| Pebble 渲染 | core/.../runners/VariableRenderer.java | VariableRenderer、render、renderOnce |
| 工作目 录 | core/.../runners/LocalWorkingDir.java | LocalWorkingDir、resolve、createTempFile |
| 文件门面 | core/.../runners/FilesService.java | FilesService、inputFiles、outputFiles |
| 日志/遮蔽 | core/.../runners/RunContextLogger.java | RunContextLogger、usedSecret、BaseAppender.replaceSecret |
| 密钥加解密 | core/.../runners/Secret.java | Secret、decrypt、encrypt |
| 插件入口 | core/.../models/tasks/RunnableTask.java | RunnableTask、run |
| MetaStore 门面 | worker/.../stores/GrpcWorkerFlowMetaStore.java | GrpcWorkerFlowMetaStore(@Replaces DefaultFlowMetaStore) |