Gas Town Scheduler 架构:配置驱动的 polecat 容量调度与延迟派发机制
2026/9/13 16:00:29 网站建设 项目流程

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 runinternal/daemon/daemon.go中的dispatchScheduledWork()使用exec.CommandContext5 分钟超时执行命令,并注入环境变量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:

字段类型描述
versionintSchema 版本(当前为 1)
work_bead_idstring被调度的实际工作 bead
target_rigstring目标 rig 名称
formulastring派发时应用的 formula(如mol-polecat-work
argsstring给执行器的自然语言指令
varsstring换行分隔的 formula 变量(key=value
enqueued_atRFC3339调度时间戳
mergestring合并策略:directmrlocal
convoystringConvoy bead ID(自动创建 convoy 后设置)
base_branchstring覆盖 polecat worktree 的基础分支
resume_branchstring恢复已有分支(与base_branch互斥)
no_mergebool完成时跳过 merge queue
review_onlybool仅评审模式:评估并汇报,不 merge/commit/push
accountstringClaude Code 账号句柄
agentstring智能体/运行时覆盖(如geminicodex
hook_raw_beadbool不带默认 formula 直接 hook
ownedbool调用方管理的 convoy 生命周期
modestring执行模式:ralph(每一步使用全新上下文)
dispatch_failuresint连续失败计数(熔断器)
last_failurestring最近一次派发错误信息

其中dispatch_failureslast_failure是熔断器的核心计数器,由recordDispatchFailure()在每次派发失败时更新(见下文"熔断器"一节)。


Bead 状态机

一个 sling context 的状态迁移如下:

+------------------+ | | v | +----------+ dispatch ok +--------+ | schedule | CONTEXT | ----------------> | CLOSED | | --------> | OPEN | | (done) | | +----------+ +--------+ | | | +-- 3 failures --> CLOSED (circuit-broken) | +-- gt scheduler clear --> CLOSED (cleared)
状态表示触发条件
SCHEDULEDOpen sling context beadscheduleBead()
DISPATCHEDClosed sling context(reason: "dispatched")dispatchSingleBead()成功
CIRCUIT-BROKENClosed sling context(reason: "circuit-broken")dispatch_failures >= 3
CLEAREDClosed 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):

  1. shouldDeferDispatch()——检查scheduler.max_polecats配置
  2. 批量(3+ 参数,最后一个是 rig)——runBatchSchedule()runBatchSling()
  3. --on标志已设置——formula-on-bead 模式
  4. 2 个参数且最后一个是 rig——scheduleBead()或内联派发
  5. 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()按顺序执行以下步骤:

  1. 校验bead 存在、rig 存在
  2. 跨 rig 守卫——若 bead 前缀与目标 rig 不匹配则拒绝(除非--force
  3. 幂等性——若该工作 bead 已存在 open sling context 则跳过
  4. 状态守卫——若 bead 处于 hooked/in_progress 则拒绝(除非--force
  5. 校验 formula——确认 formula 存在(轻量,无副作用)
  6. 烹饪 formula——bd cook在 daemon 派发前捕获坏的 protos
  7. 构建 context 字段——SlingContextFields结构体携带所有 sling 参数
  8. 创建 sling context——bd create --ephemeral+bd dep add --type=tracks(原子操作)
  9. 自动 convoy——若未被跟踪则创建 convoy,并将 convoy ID 存入 context 字段
  10. 记录事件——为 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:messagegt:handoffgt:merge-request标签的 beads 是智能体间通信工件,绝不能交给 polecat 派发capacity.IsMessagingBead()FilterMessagingBeads()(定义在 internal/scheduler/capacity/pipeline.go)在PlanDispatch的容量计算之前就做防御性剔除,readySlingContextsFromAssessments()在查询端也做同样的检查。

dispatchSingleBead

大幅简化——context 字段已经解析完毕:

  1. ReconstructFromContext(b.Context)DispatchParams,其中BeadID = b.WorkBeadID
  2. 调用executeSling(params)——就这些

派发后的清理由回调处理:

  • OnSuccessCloseSlingContext(b.ID, "dispatched")
  • OnFailure:递增dispatch_failures、更新 context bead、若熔断则关闭

ReconstructFromContext()的实现(internal/scheduler/capacity/pipeline.go)把 JSON 字段还原为DispatchParams,其中vars字符串按换行拆分回[]stringdispatchSingleBead()随后把这些参数组装成SlingParams调用executeSling(),并设置FormulaFailFatal: trueNoConvoy: trueNoBoot: true以及CallerContext: "scheduler-dispatch"——保证调度器派发时的行为一致性与幂等性。


容量管理(Capacity Management)

配置项

类型默认值描述
scheduler.max_polecats*int-1最大并发 polecats(-1=直接,0=禁用,N=延迟)
scheduler.batch_size*int1每次心跳 tick 派发的 beads 数
scheduler.spawn_delaystring"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。熔断逻辑在派发管线中有三重防线

  1. cleanupStaleContexts()在派发周期开始前就会关闭已熔断的 context(reason: "circuit-broken")以及无效 context("invalid-context")和工作 bead 已 stale 的 context("stale-work-bead",如 hooked/closed/tombstone;in_progress 有意排除——工作 bead 正在被积极处理,bd ready不会返回它,派发查询已天然防止重复派发)。
  2. assessScheduledContexts()在收集候选时跳过DispatchFailures >= maxDispatchFailures的 context。
  3. 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,包含PausedPausedByPausedAtLastDispatchAtLastDispatchCount字段。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输出pausedpaused_byqueued_totalqueued_readyactive_polecatscapacitylast_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_polecatsbatch_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.lockTryLock()。daemon 模式下拿不到锁则静默返回 0(下个周期再试),手动模式则报错"dispatch already in progress"。
  • 派发循环的不变量校验:若计划有ToDispatchDispatched == 0 && Failed == 0,直接返回scheduler dispatch invariant violation错误——这是对派发逻辑正确性的运行时断言。

代码布局(Code Layout)

路径用途
internal/scheduler/capacity/config.goSchedulerConfig类型、默认值、IsDeferred()
internal/scheduler/capacity/pipeline.goPendingBeadSlingContextFieldsPlanDispatch()ReconstructFromContext()
internal/scheduler/capacity/dispatch.goDispatchCycle类型——通用派发编排器
internal/scheduler/capacity/state.goSchedulerState持久化
internal/beads/beads_sling_context.goSling context CRUD(create、find、list、close、update)
internal/cmd/sling.goCLI 入口、配置驱动的路由
internal/cmd/sling_schedule.goscheduleBead()shouldDeferDispatch()isScheduled()
internal/cmd/scheduler.gogt scheduler命令树
internal/cmd/scheduler_epic.goEpic 调度/sling 处理器
internal/cmd/scheduler_convoy.goConvoy 调度/sling 处理器
internal/cmd/capacity_dispatch.godispatchScheduledWork()、派发回调接线
internal/daemon/daemon.go心跳集成(gt scheduler run

架构上值得注意的是internal/scheduler/capacity包刻意保持为纯函数与类型(调度循环、入队、epic/convoy 解析等不纯的编排逻辑留在 cmd 层),这使得PlanDispatchFilterCircuitBrokenCircuitBreakerPolicy等核心决策函数易于单元测试。


延伸阅读

  • 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询