Gas Town Scheduler 架构:配置驱动的 polecat 容量调度与延迟派发机制
【免费下载链接】gastownGas Town - multi-agent workspace manager项目地址: https://gitcode.com/GitHub_Trending/ga/gastown
导读
Gas Town 是一个多智能体工作区管理器,其内置的 Scheduler 组件解决了批量派发 polecat(智能体执行单元)时面临的背压(back-pressure)与容量控制问题:默认情况下gt sling一次派发 N 个 beads 会同时拉起 N 个 polecat,瞬间耗尽 API 速率限制、内存与 CPU。本文基于 docs/design/scheduler.md 及仓库源码,完整讲解 Scheduler 的核心设计——通过scheduler.max_polecats一个配置项即可在「直接派发」与「延迟派发」两种模式间无感切换,并深入剖析其 sling context bead 调度状态模型、DispatchCycle派发引擎、容量计算公式、熔断器与并发安全机制。读完本文,你将掌握如何用三条命令启用容量控制、理解调度状态如何持久化在不污染工作 bead 的独立 ephemeral bead 上,以及 daemon 心跳如何驱动增量派发。
快速开始:三步启用容量控制的延迟派发
Scheduler 是配置驱动的,不需要任何 per-command 标志位。只需设置scheduler.max_polecats配置项,同一个gt sling命令就会自动适应派发模式:
# 1. 启用延迟派发(配置驱动,无需命令级标志) gt config set scheduler.max_polecats 5 # 2. 通过 gt sling 调度工作(当 max_polecats > 0 时自动延迟) gt sling gt-abc gastown # 单个任务 bead gt sling gt-abc gt-def gt-ghi gastown # 批量任务 beads gt sling hq-cv-abc # Convoy(调度所有被跟踪的 issue) gt sling gt-epic-123 # Epic(调度所有子任务) # 3. 查看已调度的内容 gt scheduler status gt scheduler list # 4. 手动派发(或让 daemon 自动派发) gt scheduler run gt scheduler run --dry-run # 先预览派发模式(Dispatch Modes)
scheduler.max_polecats配置值完全决定派发行为:
| 值 | 模式 | 行为 |
|---|---|---|
-1(默认) | 直接派发 | gt sling立即派发,近零开销 |
0 | 直接派发 | 与-1相同——gt sling立即派发 |
N > 0 | 延迟派发 | gt sling创建 sling context bead,由 daemon 派发 |
这一语义在源码中有直接对应:internal/scheduler/capacity/config.go中的SchedulerConfig.IsDeferred()方法返回c.GetMaxPolecats() > 0,而GetMaxPolecats()在配置缺失时回退到默认值-1(直接派发)。也就是说,没有配置 = 直接派发,老用户的既有工作流完全不受影响。
常用 CLI 一览
| 命令 | 描述 |
|---|---|
gt sling <bead> <rig> | Sling bead(直接或延迟,按配置) |
gt sling <bead>... <rig> | 批量 sling/调度多个 beads |
gt sling <convoy-id> | Sling/调度 convoy 中所有被跟踪的 issues |
gt sling <epic-id> | Sling/调度 epic 的所有子任务 |
gt scheduler status | 显示调度器状态与容量 |
gt scheduler list | 按 rig 列出所有已调度的 beads |
gt scheduler run | 手动触发派发 |
gt scheduler pause | 全镇暂停所有派发 |
gt scheduler resume | 恢复派发 |
gt scheduler clear | 从调度器中移除 beads |
最小示例
gt config set scheduler.max_polecats 5 gt sling gt-abc gastown # 延迟:创建 sling context bead gt scheduler status # "Queued: 1 total, 1 ready" gt scheduler run # 派发 -> 拉起 polecat -> 关闭 context背景:为什么需要调度器
Scheduler 解决的是批量 polecat 派发的背压与容量控制问题。
没有调度器时,sling N 个 beads 会同时拉起 N 个 polecat,耗尽 API 速率限制、内存和 CPU。调度器引入了一个"调速器"(governor):beads 进入等待状态,daemon 在遵守可配置并发上限的前提下增量派发它们。
调度器作为step 14集成进 daemon 心跳流程——在所有智能体健康检查、生命周期处理和分支清理之后执行。这确保了系统在派发新工作之前是健康的:
Daemon heartbeat (every 3 min) | +- Steps 0-13: 健康检查、智能体恢复、清理 | +- Step 14: gt scheduler run (容量控制派发) | +- flock (独占锁) +- 检查暂停状态 +- 加载配置 (max_polecats, batch_size) +- 统计活跃 polecats (tmux) +- 查询 sling contexts (bd list --label=gt:sling-context) +- 与 bd ready 关联以确定未阻塞的 beads +- DispatchCycle.Run() — plan + execute + report | +- PlanDispatch(availableCapacity, batchSize, ready) | +- 对每个计划的 bead: Execute → OnSuccess/OnFailure +- 唤醒 rig 智能体 (witness, refinery) +- 保存派发状态从源码看,daemon 端通过子进程方式调用gt scheduler run:internal/daemon/daemon.go中的dispatchScheduledWork()使用exec.CommandContext以5 分钟超时执行命令,并注入环境变量GT_DAEMON=1(标识 daemon 派发,避免与手动派发混淆)和BD_DOLT_AUTO_COMMIT=off。派发门控条件是scheduler.max_polecats > 0(延迟模式)。在internal/cmd/capacity_dispatch.go中,isDaemonDispatch()通过检查GT_DAEMON == "1"来决定遇到锁冲突或暂停时是否静默跳过(daemon 模式)还是报错(手动模式)。
Sling Context Beads:调度状态的独立载体
调度状态存储在上独立的 ephemeral beads上,称为sling contexts。工作 bead 永远不会被调度器修改——这是整个设计最核心的不变量。
每个 sling context bead 具备以下特征:
- 通过
bd create --ephemeral创建,带标签gt:sling-context - 有一个
tracks依赖指向工作 bead - 所有调度参数以 JSON 形式存储在 description 中
- 在派发成功、bead 被 clear、或熔断器跳闸时关闭
为什么用独立 beads?
之前的方案在工作 bead 的 description 上存储调度元数据(分隔块),并用标签(gt:queued)作为状态信号,这需要:
- 两步写入 + 回滚(先元数据后标签)
- description 净化以避免分隔符冲突
- 三步派发清理(剥离元数据 + 交换标签 + 重试)
- 自定义 key-value 格式/解析/剥离函数(约 250 行)
Sling context beads 消除了以上所有复杂性:
- 单一原子创建——
bd create --ephemeral是一次操作 - JSON 格式——
json.Marshal/json.Unmarshal替代自定义解析器 - 工作 bead 保持原样——无 description 变更、无标签操作
- 清晰的生命周期——open context = 已调度,closed context = 已完成
源码证据在 internal/beads/beads_sling_context.go:CreateSlingContext()一次调用完成bd create --json --ephemeral --type=task --labels=gt:sling-context,随后追加dep add --type=tracks依赖(此步骤非致命——即使依赖添加失败,context bead 仍然创建成功);CloseSlingContext()对"already closed"错误做幂等抑制,保证重试安全。
Context 字段(JSON)
以下字段结构定义在 internal/scheduler/capacity/pipeline.go 的SlingContextFields结构体中,序列化为 context bead 的 description:
| 字段 | 类型 | 描述 |
|---|---|---|
version | int | Schema 版本(当前为 1) |
work_bead_id | string | 被调度的实际工作 bead |
target_rig | string | 目标 rig 名称 |
formula | string | 派发时应用的 formula(如mol-polecat-work) |
args | string | 给执行器的自然语言指令 |
vars | string | 换行分隔的 formula 变量(key=value) |
enqueued_at | RFC3339 | 调度时间戳 |
merge | string | 合并策略:direct、mr、local |
convoy | string | Convoy bead ID(自动创建 convoy 后设置) |
base_branch | string | 覆盖 polecat worktree 的基础分支 |
resume_branch | string | 恢复已有分支(与base_branch互斥) |
no_merge | bool | 完成时跳过 merge queue |
review_only | bool | 仅评审模式:评估并汇报,不 merge/commit/push |
account | string | Claude Code 账号句柄 |
agent | string | 智能体/运行时覆盖(如gemini、codex) |
hook_raw_bead | bool | 不带默认 formula 直接 hook |
owned | bool | 调用方管理的 convoy 生命周期 |
mode | string | 执行模式:ralph(每一步使用全新上下文) |
dispatch_failures | int | 连续失败计数(熔断器) |
last_failure | string | 最近一次派发错误信息 |
其中dispatch_failures与last_failure是熔断器的核心计数器,由recordDispatchFailure()在每次派发失败时更新(见下文"熔断器"一节)。
Bead 状态机
一个 sling context 的状态迁移如下:
+------------------+ | | v | +----------+ dispatch ok +--------+ | schedule | CONTEXT | ----------------> | CLOSED | | --------> | OPEN | | (done) | | +----------+ +--------+ | | | +-- 3 failures --> CLOSED (circuit-broken) | +-- gt scheduler clear --> CLOSED (cleared)| 状态 | 表示 | 触发条件 |
|---|---|---|
| SCHEDULED | Open sling context bead | scheduleBead() |
| DISPATCHED | Closed sling context(reason: "dispatched") | dispatchSingleBead()成功 |
| CIRCUIT-BROKEN | Closed sling context(reason: "circuit-broken") | dispatch_failures >= 3 |
| CLEARED | Closed sling context(reason: "cleared") | gt scheduler clear |
关键不变量:工作 bead 永不被调度器修改。所有状态都存在于 sling context bead 上。
入口点(Entry Points)
CLI 入口点
gt sling从配置和 ID 类型自动检测派发模式:
| 命令 | 直接模式(max_polecats=-1) | 延迟模式(max_polecats>0) |
|---|---|---|
gt sling <bead> <rig> | 立即派发 | 调度以便稍后派发 |
gt sling <bead>... <rig> | 批量立即派发 | 批量调度 |
gt sling <epic-id> | runEpicSlingByID()——派发所有子任务 | runEpicScheduleByID()——调度所有子任务 |
gt sling <convoy-id> | runConvoySlingByID()——派发所有被跟踪项 | runConvoyScheduleByID()——调度所有被跟踪项 |
runSling中的检测链(对应 internal/cmd/sling.go 与 internal/cmd/sling_schedule.go):
shouldDeferDispatch()——检查scheduler.max_polecats配置- 批量(3+ 参数,最后一个是 rig)——
runBatchSchedule()或runBatchSling() --on标志已设置——formula-on-bead 模式- 2 个参数且最后一个是 rig——
scheduleBead()或内联派发 - 1 个参数,自动检测类型:epic/convoy/task
在 internal/cmd/sling_schedule.go 中,shouldDeferDispatch()的判定逻辑是:找不到 town 根目录则直接返回直接派发;town settings 中无 scheduler 配置也返回直接派发;只有当GetMaxPolecats() > 0时才返回延迟派发。注意一个防御细节:若 town settings 加载失败,会返回错误并提示"修复配置或使用gt config set scheduler.max_polecats -1"——配置损坏会阻止派发,而不是静默回退。
所有调度路径都经过 internal/cmd/sling_schedule.go 中的scheduleBead()。 所有派发都经过 internal/cmd/capacity_dispatch.go 中的dispatchScheduledWork()。
Daemon 入口点
Daemon 在每次心跳(step 14)以子进程方式调用gt scheduler run:
// internal/daemon/daemon.go func (d *Daemon) dispatchScheduledWork() { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute) defer cancel() cmd := exec.CommandContext(ctx, "gt", "scheduler", "run") cmd.Env = append(os.Environ(), "GT_DAEMON=1", "BD_DOLT_AUTO_COMMIT=off") // ... }| 属性 | 值 |
|---|---|
| 超时 | 5 分钟 |
| 环境变量 | GT_DAEMON=1(标识 daemon 派发) |
| 门控 | scheduler.max_polecats > 0(延迟模式) |
调度路径(Schedule Path)
scheduleBead()按顺序执行以下步骤:
- 校验bead 存在、rig 存在
- 跨 rig 守卫——若 bead 前缀与目标 rig 不匹配则拒绝(除非
--force) - 幂等性——若该工作 bead 已存在 open sling context 则跳过
- 状态守卫——若 bead 处于 hooked/in_progress 则拒绝(除非
--force) - 校验 formula——确认 formula 存在(轻量,无副作用)
- 烹饪 formula——
bd cook在 daemon 派发前捕获坏的 protos - 构建 context 字段——
SlingContextFields结构体携带所有 sling 参数 - 创建 sling context——
bd create --ephemeral+bd dep add --type=tracks(原子操作) - 自动 convoy——若未被跟踪则创建 convoy,并将 convoy ID 存入 context 字段
- 记录事件——为 dashboard 可见性发送 feed 事件
创建是单一原子操作——无两步写入,无需回滚。
源码层面的细节值得展开:
- 幂等检查的位置:
scheduleBead()通过beads.ResolveRepoAliasBeadsDir()解析到目标 rig 的 beads 目录,再用FindOpenSlingContext()查找已存在的 open context。找到则打印Bead %s is already scheduled (context: %s), no-op并直接返回。这一设计有一个关键历史背景(GH#3468):sling context 现在创建在目标 rig 的 beads 目录而非 HQ 的 beads 目录,这样非 HQ rig 的 witness 才能在 patrol 中发现它。 - 状态守卫的纵深防御:
scheduleBead()还会拒绝 closed/tombstone 状态的 bead(bead %s is %s (work already completed)),且这一守卫不受--force绕过——如果需要重新派发,必须先 reopen 该 bead。这是为了防止 daemon 的 stranded 扫描把已完成的跨前缀 bead 重新调度,生成幽灵 convoy。 - formula 的解析顺序(
resolveFormula()):显式--formula标志 → rig property layers(gt rig config set <rig> default_formula mol-evolve,wisp 层;--global则到 bead 层)→ rig settings 文件(workflow.default_formula)→ 硬编码回退mol-polecat-work。
派发引擎(Dispatch Engine)
DispatchCycle
派发循环是一个注入回调的通用编排器:
type DispatchCycle struct { AvailableCapacity func() (int, error) // 空闲派发槽位(0=无限) QueryPending func() ([]PendingBead, error) // 有资格派发的工作项 Execute func(PendingBead) error // 派发单个项 OnSuccess func(PendingBead) error // 派发后清理 OnFailure func(PendingBead, error) // 失败处理 BatchSize int SpawnDelay time.Duration }Run()内部调用PlanDispatch(availableCapacity, batchSize, ready)决定要派发什么,然后通过回调执行每个计划项。该类型定义在 internal/scheduler/capacity/dispatch.go。
源码级增强细节:DispatchCycle还支持可选的Validate预派发钩子——返回非 nil 错误会短路该 bead 的派发(不调用Execute,直接调用OnFailure)。这用于快速不变量检查(如跨 rig 前缀守卫),它不消耗失败配额,也不会触发昂贵的派发机制。
另一个重要实现细节是OnSuccess 的重试机制(onSuccessRetries = 2):RunPlan()中OnSuccess失败会以递增间隔(attempt+1* 500ms)重试最多 3 次;若仍失败,则该 bead不计入 Dispatched,而是作为失败处理(ErrOnSuccessFailed),防止下一周期重复派发。ErrOnSuccessFailed专门用来区分"polecat 已启动但 context 关闭失败"与"polecat 从未启动"两种情况。
派发流程
DispatchCycle.Run() | +- AvailableCapacity() → capacity = maxPolecats - activePolecats | +- QueryPending() → getReadySlingContexts(): | +- bd list --label=gt:sling-context --status=open (所有 rig DBs) | +- 解析每个 context bead description 的 SlingContextFields | +- bd ready --json --limit=0 (所有 rig DBs) → readyWorkIDs 集合 | +- 过滤:WorkBeadID 在 readyWorkIDs 中的 context beads | +- 跳过熔断的(dispatch_failures >= 阈值) | +- PlanDispatch(capacity, batchSize, ready) | +- 返回 DispatchPlan{ToDispatch, Skipped, Reason} | +- 对每个计划的 bead: +- Execute: ReconstructFromContext(fields) → executeSling(params) +- OnSuccess: CloseSlingContext(contextID, "dispatched") +- OnFailure: 递增 dispatch_failures、更新 context、必要时关闭 +- sleep(SpawnDelay)源码级增强细节:getReadySlingContexts()的实现(在 internal/cmd/capacity_dispatch.go)比文档中的伪代码更精细——它先通过assessScheduledContexts()对每个 open context 做批量评估:按enqueued_at排序(FIFO,先调度先派发),批量获取工作 bead 状态(batchFetchBeadInfoByIDs使用bd show --json按 beads 目录分组批量查询,避免对大型仓库执行 O(minutes) 的bd list --all),再通过bd blocked --json查询阻塞状态。一个 bead 被视为 ready 的条件是:工作 bead 存在、未被阻塞、状态为 open。
此外,派发管线中有两道针对消息类标签的防御过滤(引用自 gt-el4 事件):gt:message、gt:handoff、gt:merge-request标签的 beads 是智能体间通信工件,绝不能交给 polecat 派发。capacity.IsMessagingBead()与FilterMessagingBeads()(定义在 internal/scheduler/capacity/pipeline.go)在PlanDispatch的容量计算之前就做防御性剔除,readySlingContextsFromAssessments()在查询端也做同样的检查。
dispatchSingleBead
大幅简化——context 字段已经解析完毕:
ReconstructFromContext(b.Context)→DispatchParams,其中BeadID = b.WorkBeadID- 调用
executeSling(params)——就这些
派发后的清理由回调处理:
- OnSuccess:
CloseSlingContext(b.ID, "dispatched") - OnFailure:递增
dispatch_failures、更新 context bead、若熔断则关闭
ReconstructFromContext()的实现(internal/scheduler/capacity/pipeline.go)把 JSON 字段还原为DispatchParams,其中vars字符串按换行拆分回[]string。dispatchSingleBead()随后把这些参数组装成SlingParams调用executeSling(),并设置FormulaFailFatal: true、NoConvoy: true、NoBoot: true以及CallerContext: "scheduler-dispatch"——保证调度器派发时的行为一致性与幂等性。
容量管理(Capacity Management)
配置项
| 键 | 类型 | 默认值 | 描述 |
|---|---|---|---|
scheduler.max_polecats | *int | -1 | 最大并发 polecats(-1=直接,0=禁用,N=延迟) |
scheduler.batch_size | *int | 1 | 每次心跳 tick 派发的 beads 数 |
scheduler.spawn_delay | string | "0s" | 两次 spawn 之间的延迟(避免 Dolt 锁竞争) |
通过gt config set设置:
gt config set scheduler.max_polecats 5 # 启用延迟派发 gt config set scheduler.max_polecats -1 # 直接派发(默认) gt config set scheduler.batch_size 2 gt config set scheduler.spawn_delay 3s源码级细节:SchedulerConfig定义在 internal/scheduler/capacity/config.go,是一个全镇级(town-wide)设置而非 per-rig——因为 API 速率限制、内存和 CPU 是所有 rig 共享的宿主级资源。GetMaxPolecats()、GetBatchSize()、GetSpawnDelay()都在字段缺失时回退到默认值(-1、1、"0s");ParseDurationOrDefault()对非法时长字符串同样回退到 fallback。scheduler run --batch N可在运行时覆盖 batch_size(batchOverride > 0时生效)。
派发数量公式
toDispatch = min(capacity, batchSize, readyCount) 其中: capacity = maxPolecats - activePolecats(正数 = 空闲槽位数,0 或负数 = 无容量) batchSize = scheduler.batch_size(默认 1) readyCount = 工作 bead 出现在 bd ready 中的 sling context 数PlanDispatch()是这一公式的纯函数实现(internal/scheduler/capacity/pipeline.go),它会返回一个DispatchPlan{ToDispatch, Skipped, Reason},其中Reason精确标识本次限制因素:"capacity"(容量不足)、"batch"(达到批次上限)、"ready"(就绪数量不足)、"none"(无就绪 bead)。dry-run 模式下这些原因会直接展示给操作员。
活跃 Polecat 计数
活跃 polecat 通过扫描 tmux 会话并调用session.ParseSessionName()匹配角色来统计(countActivePolecats()在 internal/cmd/scheduler.go)。这会统计所有polecat——包括调度器派发的和直接 sling 的——因为 API 速率限制、内存和 CPU 是共享资源。
在派发主路径上,实际用于容量准入的是更精细的polecatCapacitySnapshotForTown(),它区分 working、recovery_blocked、reservations、reusable_idle、pending_mr 等状态,gt scheduler status会把这些细分维度完整展示出来。
熔断器(Circuit Breaker)
熔断器防止永远失败的 beads 导致无限重试循环。
| 属性 | 值 |
|---|---|
| 阈值 | maxDispatchFailures = 3 |
| 计数器 | sling context JSON 中的dispatch_failures字段 |
| 跳闸动作 | 关闭 sling context(reason: "circuit-broken") |
| 重置 | 无自动重置(需人工干预) |
流程
派发尝试失败 | +- 递增 context bead 中的 dispatch_failures +- 存储 last_failure 错误信息 | +- dispatch_failures >= 3? +- 是 -> CloseSlingContext(contextID, "circuit-broken") | (context bead 关闭,工作 bead 不受影响) +- 否 -> bead 保持已调度状态,下一周期重试源码级细节:maxDispatchFailures = 3定义在 internal/cmd/capacity_dispatch.go。recordDispatchFailure()递增计数并记录错误;达到阈值后关闭 context。熔断逻辑在派发管线中有三重防线:
cleanupStaleContexts()在派发周期开始前就会关闭已熔断的 context(reason: "circuit-broken")以及无效 context("invalid-context")和工作 bead 已 stale 的 context("stale-work-bead",如 hooked/closed/tombstone;in_progress 有意排除——工作 bead 正在被积极处理,bd ready不会返回它,派发查询已天然防止重复派发)。assessScheduledContexts()在收集候选时跳过DispatchFailures >= maxDispatchFailures的 context。capacity.FilterCircuitBroken()作为纯函数提供最终的过滤工具。
另外,PlanDispatch的失败策略在纯函数层有抽象:CircuitBreakerPolicy(maxFailures)返回"达到阈值前重试、之后隔离"的策略,NoRetryPolicy()则首次失败即隔离。
调度器控制(Scheduler Control)
Pause / Resume
暂停会全镇停止所有派发。状态存储在.runtime/scheduler-state.json。
gt scheduler pause # 设置 paused=true,记录操作者与时间戳 gt scheduler resume # 清除暂停状态写入是原子的(临时文件 + 重命名),防止并发写入者造成损坏。
源码级细节:SchedulerState结构体定义在 internal/scheduler/capacity/state.go,包含Paused、PausedBy、PausedAt、LastDispatchAt、LastDispatchCount字段。LoadState()在文件不存在时返回零值状态(有意设计:缺失 = "未暂停、从未派发"),并支持从旧的queue-state.json迁移;SaveState()采用os.CreateTemp+os.Rename的原子写入。文档强调的"fresh state on save"在 internal/cmd/capacity_dispatch.go 中落地为:派发完成后重新读取状态再写入RecordDispatch(),避免覆盖并发 pause 操作。
Clear
关闭 sling context beads,将 beads 从调度器中移除:
gt scheduler clear # 关闭 ALL sling contexts gt scheduler clear --bead gt-abc # 关闭特定 bead 的 context源码级细节:--bead变体会扫描所有rig 目录下的 contexts(因为 context 存在于目标 rig 的 beads 目录,GH#3468),并关闭该工作 bead 对应的全部context(处理并发scheduleBead竞态可能产生的重复 context)。
Status / List
gt scheduler status # 摘要:paused、queued 数量、活跃 polecats gt scheduler status --json # JSON 输出 gt scheduler list # 按目标 rig 分组的 beads,带阻塞指示符 gt scheduler list --json # JSON 输出list将 sling contexts(所有已调度项)与bd ready(未阻塞的工作 beads)对账,以标记阻塞的 beads。status --json输出paused、paused_by、queued_total、queued_ready、active_polecats、capacity、last_dispatch_at等结构化字段,便于脚本消费。
调度器与 Convoy 的集成
Convoys 和调度器是互补但不同的机制。Convoys 跟踪相关 beads 的完成情况;调度器控制派发容量。派发 convoy 工作有两条路径:
派发路径
| 路径 | 触发 | 容量控制 | 使用场景 |
|---|---|---|---|
| 直接派发 | gt sling <convoy-id>(max_polecats=-1) | 无(立即触发) | 默认模式——所有 issues 一次性派发 |
| 延迟派发 | gt sling <convoy-id>(max_polecats>0) | 有(daemon 心跳、max_polecats、batch_size) | 容量控制——批量 + 背压 |
直接派发(max_polecats=-1):gt sling <convoy-id>调用runConvoySlingByID(),通过executeSling()立即派发所有 open 的被跟踪 issues。每个 issue 的 rig 从其 bead ID 前缀自动解析。无容量控制——所有 issues 同时派发。
延迟派发(max_polecats>0):gt sling <convoy-id>调用runConvoyScheduleByID(),调度所有 open 的被跟踪 issues(创建 sling context beads)。Daemon 通过gt scheduler run增量派发,遵守max_polecats和batch_size。对于同时派发会耗尽资源的大批量任务,请使用此模式。
何时用哪种
- 小 convoy(< 5 个 issues):直接派发(默认,max_polecats=-1)
- 大批量(5+ 个 issues):设置
scheduler.max_polecats以启用容量控制派发 - Epics:同样的逻辑——
gt sling <epic-id>从配置自动解析模式
Rig 解析
gt sling <convoy-id>和gt sling <epic-id>通过beads.ExtractPrefix()+beads.GetRigNameForPrefix()从每个 bead 的 ID 前缀自动解析目标 rig。Town 根 beads(hq-*)会被跳过并给出警告,因为它们是协调工件而非可派发的工作。
detectSchedulerIDType()的类型检测顺序为:hq-cv-前缀快速路径 → bead 的IssueType(epic/convoy)→ 标签(gt:epic/gt:convoy)→ 回退为 task。注意 convoy/epic 模式下会校验不允许使用 task-only 标志(--account、--agent、--ralph、--args、--var、--merge、--base-branch、--no-convoy、--owned、--no-merge、--review-only)。
安全属性(Safety Properties)
| 属性 | 机制 |
|---|---|
| 调度幂等性 | 若工作 bead 已存在 open sling context 则跳过 |
| 工作 bead 保持原样 | 调度器从不修改工作 bead 的 description 或标签 |
| 跨 rig 守卫 | 若 bead 前缀与目标 rig 不匹配则拒绝(除非--force) |
| 派发串行化 | flock(scheduler-dispatch.lock)防止双重派发 |
| 原子调度 | 单次bd create --ephemeral——无两步写入,无回滚 |
| formula 预烹饪 | 调度时bd cook在 daemon 派发循环前捕获坏的 protos |
| 保存时读取最新状态 | 派发在保存前重新读取状态,避免覆盖并发 pause |
源码级细节:
- 跨 rig 守卫有两层。调度时:
checkCrossRigGuard()(除非--force);派发时:validatePendingBeadForDispatch()调用capacity.AcceptsPrefix()检查前缀匹配(定义于 internal/scheduler/capacity/dispatch.go)。若 rig 前缀未知(空),退化为接受(开放降级而非拒绝派发)。派发层发现跨 rig 前缀不匹配会打印告警并触发gt escalateMEDIUM 级别告警(带 1 小时防抖,防止每个心跳 tick 都刷屏)。 - 派发串行化:
dispatchScheduledWork()使用gofrs/flock对.runtime/scheduler-dispatch.lock做TryLock()。daemon 模式下拿不到锁则静默返回 0(下个周期再试),手动模式则报错"dispatch already in progress"。 - 派发循环的不变量校验:若计划有
ToDispatch但Dispatched == 0 && Failed == 0,直接返回scheduler dispatch invariant violation错误——这是对派发逻辑正确性的运行时断言。
代码布局(Code Layout)
| 路径 | 用途 |
|---|---|
| internal/scheduler/capacity/config.go | SchedulerConfig类型、默认值、IsDeferred() |
| internal/scheduler/capacity/pipeline.go | PendingBead、SlingContextFields、PlanDispatch()、ReconstructFromContext() |
| internal/scheduler/capacity/dispatch.go | DispatchCycle类型——通用派发编排器 |
| internal/scheduler/capacity/state.go | SchedulerState持久化 |
| internal/beads/beads_sling_context.go | Sling context CRUD(create、find、list、close、update) |
| internal/cmd/sling.go | CLI 入口、配置驱动的路由 |
| internal/cmd/sling_schedule.go | scheduleBead()、shouldDeferDispatch()、isScheduled() |
| internal/cmd/scheduler.go | gt scheduler命令树 |
| internal/cmd/scheduler_epic.go | Epic 调度/sling 处理器 |
| internal/cmd/scheduler_convoy.go | Convoy 调度/sling 处理器 |
| internal/cmd/capacity_dispatch.go | dispatchScheduledWork()、派发回调接线 |
| internal/daemon/daemon.go | 心跳集成(gt scheduler run) |
架构上值得注意的是internal/scheduler/capacity包刻意保持为纯函数与类型(调度循环、入队、epic/convoy 解析等不纯的编排逻辑留在 cmd 层),这使得PlanDispatch、FilterCircuitBroken、CircuitBreakerPolicy等核心决策函数易于单元测试。
延伸阅读
- Convoys——Convoy 跟踪、调度时的自动 convoy 创建
- Property Layers——调度器标签使用的 labels-as-state 模式(见 Operational State Events 章节)
- 调度派发作为 daemon 心跳 step 14 运行,其上游是 daemon 心跳的完整健康检查链(docs/design/ 目录下的 watchdog 相关设计文档)
以上所有源码路径与配置示例均以当前仓库为准。若需在本地复现本文示例,请先在已初始化的 Gas Town town 目录内执行gt config set scheduler.max_polecats N后运行gt scheduler status观察状态变化,并使用gt scheduler run --dry-run在真实派发前预览计划。
【免费下载链接】gastownGas Town - multi-agent workspace manager项目地址: https://gitcode.com/GitHub_Trending/ga/gastown
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考