Archon 架构深度剖析:从 Slack 消息到工作流执行的全链路数据流
【免费下载链接】ArchonThe first open-source harness builder for AI coding. Make AI coding deterministic and repeatable.项目地址: https://gitcode.com/GitHub_Trending/archon3/Archon
本文基于仓库文档 .claude/docs/architecture-deep-dive.md 展开,结合 Archon 各子包源码(core / workflows / isolation / server / web / adapters)进行交叉印证。全文以端到端的数据流追踪为主线,覆盖消息路由、工作流执行、隔离解析、Session 状态机、数据库抽象、配置加载、SSE 实时推送等核心链路,并附关键文件索引,可作为跨系统排障与二次开发的架构导览。
Archon 是一套面向 AI 编码场景的开源 harness 构建器,目标是让 AI 编码变得可确定、可复现。它并不只是"接一个 Slack/Web 聊天机器人",而是把消息入口、编排 Agent、工作流执行器、worktree 隔离、数据库、配置与实时 UI 串成一条可追踪的流水线。本文按数据流动的顺序逐层拆解:一条 Slack 消息进来后如何被路由、如何触发工作流、如何在工作树中隔离执行、如何落库、如何把进度推送到 Web UI,以及贯穿其中的横切模式。
1. 消息流转:路由代理架构
核心结论先行:编排器(orchestrator)本质是一个路由代理(routing agent)——绝大多数消息都会经过一次 AI 调用,由模型决定如何处理,而不是走命令分发器。这与传统"命令分发"式机器人有本质区别:即使输入是未知的斜杠命令,也交给 AI 解读。
1.1 入口:Slack 事件到锁管理
以 Slack 为例,完整链路如下(源码路径见 packages/adapters/src/chat/slack/adapter.ts):
Slack event → SlackAdapter.start() 注册 app_mention + message 处理器 → 授权检查:isSlackUserAuthorized()(auth.ts:27) → void this.messageHandler(event) —— fire-and-forget(adapter.ts:416/457) → lockManager.acquireLock(conversationId, handler)(conversation-lock.ts:59) → handleMessage(platform, conversationId, text) → db.getOrCreateConversation() → inheritThreadContext() —— 若为子线程,复制父线程的 codebase/cwd → generateAndSetTitle() —— 非斜杠消息异步生成标题几个值得注意的入口细节:
- 双事件入口:
app_mention处理被 @提及的消息(adapter.ts:392);message事件只处理channel_type === 'im'的私聊,并跳过 bot 自身消息以避免死循环(adapter.ts:421-433)。 - 静默拒绝未授权用户:授权检查失败时仅记录掩码日志(如
U123***),不回复任何内容,避免向未授权用户暴露 bot 存在(adapter.ts:395-399)。 - 斜杠命令也走同一链路:
/archon与/archon-workflow两个 Slack 斜杠命令(adapter.ts:461-468)会先发送一条可见的 "seed" 消息作为线程根,再把转换后的文本注入同一messageHandler(adapter.ts:524-581)。其中/archon connect github被内联处理,走 GitHub 设备流。 - Fire-and-forget 模式:
void this.messageHandler(...)表示调用方不阻塞等待;错误由消息处理器内部兜底。
1.2 并发控制:Conversation Lock Manager
锁管理器(packages/core/src/utils/conversation-lock.ts)保证两点:
- 全局并发上限:默认
maxConcurrent = 10(构造器第 46 行),同时最多处理 10 个会话; - 单会话串行:同一 conversation 内的消息按顺序处理,后到消息进入队列。
acquireLock是非阻塞的:若会话已活跃,返回{ status: 'queued-conversation' };若已达全局容量,返回{ status: 'queued-capacity' };否则立即执行并返回{ status: 'started' }(conversation-lock.ts:59-110)。handler 的 Promise 会先写入 Map 再 await,以消除竞态;完成时再触发队列与全局队列的推进(conversation-lock.ts:82-104)。
1.3 决策分叉:命令还是 AI?
消息进入handleMessage后分成两条路径(packages/core/src/orchestrator/orchestrator-agent.ts):
路径 A —— 确定性命令处理(约 5 个白名单命令):
IF message 以 '/' 开头 且 命令 ∈ [help, status, reset, workflow, register-project]: → commandHandler.handleCommand()(确定性分发,不经过 AI) → 若结果为 workflow → handleWorkflowRunCommand() → dispatchOrchestratorWorkflow() → 直接返回响应路径 B —— AI 路由(其余所有消息,包括未知斜杠命令):
→ codebaseDb.listCodebases() + discoverAllWorkflows() → buildFullPrompt()(prompt-builder.ts) → 会话挂接 codebase 时 → buildProjectScopedPrompt() → 否则 → buildOrchestratorPrompt()(罗列所有已注册项目) → Prompt 内含:已注册项目、已发现工作流、/invoke-workflow 格式说明 → sessionDb.getActiveSession();若无 → transitionSession('first-message') → getAgentProvider(conversation.ai_assistant_type) → cwd = getArchonWorkspacesPath() → 按 getStreamingMode() 选择 handleBatchMode() / handleStreamMode() AI 回复自然语言 ± 结构化命令: → filterToolIndicators(assistantMessages) —— 剥离 emoji 前缀的工具噪音 → parseOrchestratorCommands() → 命中 /invoke-workflow → dispatchOrchestratorWorkflow() → 命中 /register-project → handleRegisterProject() → 否则 → 剩余文本经 platform.sendMessage() 发给用户值得注意:Prompt 中出现的/invoke-workflow由 AI 决定是否发出,也就是说工作流派发的"最后一公里"决定权在模型手里;parseOrchestratorCommands则负责从回复中提取这些结构化指令。源码中还有针对模型输出格式的防御逻辑,例如normalizeCommandText(orchestrator-agent.ts:401-403)会剥离**\/register-project ...**这类被 markdown 加粗污染的命令行,保证isCommandFullyParsed能正确识别。
1.4 关键决策点小结
| 决策点 | 行为 |
|---|---|
getStreamingMode() | Slack 返回'batch',Web 返回'stream' |
buildFullPrompt() | 有 codebase 用项目级 prompt,否则用全局 orchestrator prompt |
parseOrchestratorCommands() | 由 AI 决定派发工作流还是纯对话回复 |
| Session resume | 将session.assistant_session_id传入 SDK 的options.resume |
| 确定性命令数量 | 仅 5 个;其余一切(含斜杠命令)均由 AI 路由 |
2. 工作流执行:/workflow run archon-fix-github-issue #42
2.1 派发链路
用户消息以 /workflow 开头 → commandHandler.handleCommand()(orchestrator-agent.ts:422) → discoverWorkflowsWithConfig() 按名称找到工作流(workflows/src/loader.ts) → 返回 CommandResult,result.workflow = { definition, args } → handleWorkflowRunCommand()(orchestrator-agent.ts:888) → dispatchOrchestratorWorkflow()(orchestrator-agent.ts:192) → validateAndResolveIsolation() —— 见第 3 节 → 非 Web 渠道:executeWorkflow() 直接执行(orchestrator-agent.ts:249) → Web 渠道:dispatchBackgroundWorkflow() → 独立 worker 会话 + fire-and-forget(orchestrator.ts:336)Web 与 Slack/CLI 的关键差异就在这里:Web 端为了不阻塞用户页面,把工作流放到后台 worker 会话执行,执行事件通过事件桥转发回父会话的 SSE 流(见第 7.3 节)。
2.2executeWorkflow()内部(executor.ts)
→ deps.store.createWorkflowRun() —— 创建 DB 运行记录 → getWorkflowEventEmitter().registerRun(runId, conversationId) → 从配置解析 provider/model → 创建 artifactsDir 与 logDir IF isDagWorkflow: → executeDagWorkflow()(dag-executor.ts) → buildTopologicalLayers() —— Kahn 算法拓扑分层 → 每层 Promise.allSettled(nodes) 并行执行 → 每节点:checkTriggerRule() → evaluateCondition(when) → bash 节点:execFileAsync('bash', ['-c', script]) → AI 节点:resolveNodeProviderAndModel() → aiClient.sendQuery() → 输出写入 nodeOutputs map,供 $nodeId.output 引用 IF isLoopWorkflow: → for i = 1..max_iterations: → substituteWorkflowVariables(prompt) → aiClient.sendQuery() → detectCompletionSignal(output, until) —— 命中停止条件则 break IF isStepWorkflow: → for each step: → SingleStep:executeStepInternal() → loadCommandPrompt(cwd, commandName) —— 先搜仓库再回落 bundled defaults → substituteWorkflowVariables() —— 替换 $ARGUMENTS、$ARTIFACTS_DIR 等变量 → withIdleTimeout(aiClient.sendQuery(), idleTimeout) → 流式或批量输出到平台 → ParallelBlock:Promise.all(executeStepInternal per sub-step)三种执行模型的本质差异:
- DAG:节点间显式声明依赖,由 Kahn 算法分层,同层并行,支持
when条件与$nodeId.output引用; - Loop:迭代执行并检测"完成信号"(
until),适合"反复修改直到满足条件"的循环任务; - Step:顺序执行步骤,支持单步与并行块,变量替换面向命令文件模板。
2.3 事件发射
每一步/节点通过WorkflowEventEmitter发射step_started、step_completed、node_started等事件,再由WorkflowEventBridge转发为 SSE 事件推送到 Web UI(dag_node、workflow_step等),驱动前端的工作流进度卡片实时刷新。
3. 隔离解析:7 步 Worktree 算法
3.1 解析入口
validateAndResolveIsolation()(orchestrator.ts:108) → IsolationResolver.resolve(request)(isolation/src/resolver.ts:100)IsolationResolver按严格优先级执行 7 步探测,尽量复用已有环境,避免无谓地新建 worktree(packages/isolation/src/resolver.ts):
| 步骤 | 探测内容 | 结果 |
|---|---|---|
| 1 | store.getById(envId)+worktreeExists()检查既有环境 | 有效 →{ status: 'resolved', method: 'existing' };过期 →markDestroyedBestEffort()→{ status: 'stale_cleaned' }并让调用方重试 |
| 2 | 无 codebase | { status: 'none', cwd: '/workspace' } |
| 3 | 工作流复用:store.findActiveByWorkflow(codebaseId, workflowType, workflowId) | 有效 →{ method: 'workflow_reuse' } |
| 4 | 关联 issue:遍历hints.linkedIssues,找活动的issue环境 | 命中 →{ method: 'linked_issue_reuse' } |
| 5 | PR 分支收养:findWorktreeByBranch(canonicalPath, prBranch) | 命中 →store.create({ adopted: true })→{ method: 'branch_adoption' } |
| 6 | 数量上限:store.countActiveByCodebase()vsmaxWorktrees(25) | 触顶 →cleanup.makeRoom()清理后重查,仍满则拒绝 |
| 7 | 新建:provider.create(isolationRequest)→store.create() | store.create()失败时销毁孤儿 worktree 并重抛 |
3.2WorktreeProvider.create()内部(worktree.ts:56)
→ generateBranchName(request) —— 按场景生成:issue-N、thread-{hash}、task-{slug} 等 → getWorktreePath() —— ~/.archon/workspaces/{owner}/{repo}/worktrees/{branch} → findExisting() —— 检查路径或 PR 分支是否可收养 → syncWorkspaceBeforeCreate() —— git fetch origin {baseBranch} → git worktree add {path} -b {branch} origin/{baseBranch} → copyConfiguredFiles() —— 复制 .archon/ 与 config.worktree.copyFiles 中声明的文件这里的"收养"(adoption)机制非常实用:如果目标分支上已经存在一个旧 worktree,Archon 不会重复创建,而是直接接管它,从而保留现场、节省磁盘并避免冲突。
4. Session 生命周期:状态机
最重要的设计原则:Session 迁移是"不可变"的——已有 session 永远不会被修改,只会被停用并替换。这样每段会话的历史(assistant_session_id、上下文)都保留在只读记录中,可随时回溯。
4.1 典型迁移流程
首条消息 → transitionSession('first-message') → INSERT 新 session(parent_session_id = null) → assistant_session_id = null(尚无 SDK 会话) AI 调用完成 → tryPersistSessionId(session.id, sdkSessionId) → UPDATE assistant_session_id,供下一条消息 resume 下一条消息 → getActiveSession() 返回既有 session → sendQuery(..., session.assistant_session_id) —— SDK 自动续接上下文 /reset → transitionSession('reset-requested') → 停用当前 session(ended_reason = 'reset-requested') → 不立即创建新 session → 下一条消息触发 'first-message' → 创建新 session Plan → Execute 迁移: → detectPlanToExecuteTransition() 检测 commandName === 'execute' && lastCommand === 'plan-feature' → transitionSession('plan-to-execute') —— 唯一会立即创建新 session 的触发器 → 旧 session 停用 + 新 session 创建,在同一个 DB 事务内原子完成4.2 TransitionTrigger 枚举
文档给出的完整触发器集合如下:
'first-message'、'plan-to-execute'、'isolation-changed'、'codebase-changed'、 'codebase-cloned'、'cwd-changed'、'reset-requested'、'context-reset'、 'repo-removed'、'worktree-removed'、'conversation-closed'当前源码中的枚举(packages/core/src/state/session-transitions.ts)实现了其中的核心子集,并明确划分为三种行为类别:
const TRIGGER_BEHAVIOR: Record<TransitionTrigger, 'creates' | 'deactivates' | 'none'> = { 'first-message': 'none', // 没有既有 session 可停用 'plan-to-execute': 'creates', // 唯一"停用 + 立即创建"的场景 'isolation-changed': 'deactivates', 'project-changed': 'deactivates', 'reset-requested': 'deactivates', 'worktree-removed': 'deactivates', 'conversation-closed': 'deactivates', };'creates':停用当前 session并立即创建新 session;'deactivates':仅停用当前 session,由下一条消息触发新 session;'none':不执行任何操作。
该 Record 类型保证了编译期穷尽性——新增触发器若未归类会直接触发 TypeScript 编译错误。这一设计在源码注释中被明确为"单一事实来源"(single source of truth)。
4.3 审计追踪
getSessionChain(sessionId)通过递归 CTE 沿parent_session_id链接回溯整个会话链,把"这次对话经历过哪些上下文切换"完整还原出来——这既是调试会话上下文问题的利器,也是审计"为什么 AI 在某个节点丢失上下文"的依据。
5. 数据库层:IDatabase 抽象
5.1 自动检测(connection.ts:30-46)
DATABASE_URL 已设置 → PostgresAdapter(pg.Pool,max: 10) 否则 → SqliteAdapter(bun:sqlite,WAL 模式,busy_timeout: 5000)源码中的getDatabaseType()(packages/core/src/db/connection.ts)同样是单一判断:process.env.DATABASE_URL ? 'postgresql' : 'sqlite'。也就是说,零配置默认使用 SQLite,设置DATABASE_URL即无缝切换 PostgreSQL。
5.2 查询流与方言适配
- PostgreSQL:
$1、$2占位符原生可用; - SQLite:
convertPlaceholders()把$N替换为?并重排参数,同时剥离::jsonb类型转换——这样同一份 SQL 可以跑在两个引擎上。
5.3 命名空间导出模式
import * as conversationDb from '@archon/core/db/conversations'; import * as sessionDb from '@archon/core/db/sessions'; await conversationDb.getOrCreateConversation(platformType, conversationId); await sessionDb.transitionSession(conversationId, trigger, options);每个 db 模块以命名空间方式整体导出,调用方按领域(conversations / sessions / messages / users ...)聚合引用,避免大而全的单一 DB 对象。
5.4 方言差异对照表
| Feature | SQLite | PostgreSQL |
|---|---|---|
now() | datetime('now') | NOW() |
jsonMerge(col, $N) | json_patch(col, $N) | col \|\| $N::jsonb |
| UUID | crypto.randomUUID() | gen_random_uuid() |
6. 配置加载:4 层合并
配置不是单一文件,而是按优先级从低到高合并 4 层(packages/core/src/config/config-loader.ts):
Layer 1: 代码默认值(config-loader.ts) → botName: 'Archon'、assistant: 'claude'、concurrency.maxConversations: 10 Layer 2: 全局配置(~/.archon/config.yaml) → loadGlobalConfig() —— 首次加载后缓存 → 覆盖项:botName、defaultAssistant、assistants.*、流式模式 Layer 3: 仓库配置({repoPath}/.archon/config.yaml) → loadRepoConfig() —— 每次读取不缓存 → 覆盖项:assistant、assistants.*、commands.folder、defaults.*、worktree.baseBranch Layer 4: 环境变量(优先级最高) → BOT_DISPLAY_NAME、DEFAULT_AI_ASSISTANT → TELEGRAM_STREAMING_MODE、DISCORD_STREAMING_MODE、SLACK_STREAMING_MODE → MAX_CONCURRENT_CONVERSATIONS两个值得注意的工程决策:
- 全局配置缓存,仓库配置不缓存:全局配置变化频率低、影响面大,首次加载后缓存(
loadGlobalConfig(forceReload)支持强制刷新);仓库配置则"读新鲜"(loadRepoConfig(repoPath)),因为仓库的.archon/config.yaml可能被工作流动态修改,需要即时生效。 - 环境变量穿透:
MAX_CONCURRENT_CONVERSATIONS会覆盖到concurrency.maxConversations(源码 config-loader.ts:506 解析该环境变量),与锁管理器的maxConcurrent联动。
工作流模型解析优先级
当工作流执行需要解析模型时,按以下顺序回落:
- 节点级
model(DAG 模式,per-node 声明); - 工作流级
model(YAML 顶层声明); - 配置
assistants.{provider}.model; - SDK 默认值。
在聊天场景下,源码 orchestrator-agent.ts 的resolveChatModelRequest还引入了更高优先级的 per-user 偏好(default_model+default_provider匹配时才生效),并以"tier"(层级别名)机制做整体回落——这说明模型解析在聊天与工作流两条路径上是分别实现的,工作流路径保持large层级语义。
7. Web UI 数据流:React → SSE → Server
Web 端采用REST(TanStack Query v5)承载静态数据 + SSE 承载实时增量的双通道架构。
7.1 REST 数据(TanStack Query v5)
React 组件 → useQuery({ queryKey, queryFn }) → apiClient.listConversations() —— fetch('/api/conversations') → Server:Hono 路由处理器 → DB 查询 → JSON 响应 → TanStack Query 负责缓存、轮询、失效7.2 SSE 实时流
React:useSSE(conversationId)(web/src/hooks/useSSE.ts) → new EventSource(`${SSE_BASE_URL}/api/stream/${conversationId}`) → Server:streamSSE(c, async (stream) => { transport.registerStream(conversationId, stream) stream.onAbort(() => transport.removeStream(...)) }) 事件流: AI client 产出内容 → WebAdapter.sendMessage() → persistence.appendText() —— 先缓冲,稍后落库 → transport.emit(conversationId, { type: 'text', content }) → stream.writeSSE({ data: JSON.stringify(event) }) 客户端接收: → eventSource.onmessage → parseSSEEvent() → switch(data.type): 'text' → 50ms debounce 缓冲 → handlers.onText() 'tool_call' → flush text → handlers.onToolCall() 'tool_result' → flush text → handlers.onToolResult() 'conversation_lock' → handlers.onLockChange() 'workflow_step' → handlers.onWorkflowStep() 'dag_node' → handlers.onDagNode() 'retract' → 清空缓冲 → handlers.onRetract()retract(撤回)事件的存在说明前端渲染与 AI 流式输出之间存在一致性协调机制:当模型撤回之前的内容时,客户端会清空 debounce 缓冲并触发重渲染,避免展示脏文本。
7.3 工作流进度(后台工作流)
工作流执行器发事件 → WorkflowEventEmitter 单例 → WorkflowEventBridge 订阅 → mapWorkflowEvent() → 后台工作流:bridgeWorkerEvents(workerConvId, parentConvId) → 把 worker 会话的事件路由到父会话的 SSE 流 → transport.emitWorkflowEvent(parentConvId, sseEvent) → SSE → React → WorkflowProgressCard 更新这条链路的精巧之处:后台工作流运行在独立的 worker 会话中(与前端页面解耦),但事件桥会把 worker 的进度事件透明地转发到父会话流上,用户看到的进度体验与前台执行几乎一致。
7.4 重连宽限期
SSETransport.removeStream()会调度RECONNECT_GRACE_MS = 5000ms后的清理。若客户端在 5 秒内重连(典型场景:浏览器路由跳转导致的 EventSource 断开),registerStream()会取消清理定时器——持久化缓冲状态得以保留,用户不会因短暂跳转而丢失正在流式输出的内容。
8. 横切模式(Cross-Cutting Patterns)
8.1 Lazy Logger(延迟日志器)
每个模块都延迟创建 logger,避免测试 mock 时机问题:
let cachedLog: ReturnType<typeof createLogger> | undefined; function getLog() { return (cachedLog ??= createLogger('module')); }该模式在 conversation-lock.ts、session-transitions.ts 等处反复出现——测试可以在createLogger被首次调用前注入 mock。
8.2execFileAsync(而非exec)
所有 git 子进程调用统一走 packages/git/src/exec.ts,用参数数组而非 shell 字符串拼接,从根上避免 shell 注入,并提供一致的超时处理。
8.3 结构化事件旁路(Structured Event Side-Channel)
IPlatformAdapter.sendStructuredEvent?()是可选方法,仅WebAdapter实现。编排器与执行器在调用前都会检查if (platform.sendStructuredEvent)。作用是把 SDK 的原始 tool call 对象与格式化文本分开,经独立通道推给 SSE——前端因此可以拿到结构化的工具调用数据渲染专用 UI,而不会被 markdown 格式化污染。
8.4isWebAdapter()类型守卫
将IPlatformAdapter收窄为WebAdapter,以便安全调用 Web 专属方法:setConversationDbId()、setupEventBridge()、emitRetract()。这是 TypeScript 类型守卫在跨平台抽象中的典型用法。
9. 关键文件索引
| 流程 | 关键文件 |
|---|---|
| 消息入口 | packages/adapters/src/chat/slack/adapter.ts、packages/server/src/index.ts |
| 编排 | packages/core/src/orchestrator/orchestrator-agent.ts、packages/core/src/orchestrator/orchestrator.ts |
| 锁管理 | packages/core/src/utils/conversation-lock.ts |
| AI Provider | packages/providers/src/claude/index.ts、packages/providers/src/registry.ts |
| 命令处理 | packages/core/src/handlers/command-handler.ts |
| Session | packages/core/src/db/sessions.ts、packages/core/src/state/session-transitions.ts |
| 工作流 | packages/workflows/src/executor.ts、packages/workflows/src/dag-executor.ts、packages/workflows/src/loader.ts |
| 隔离 | packages/isolation/src/resolver.ts、packages/isolation/src/providers/worktree.ts |
| 数据库 | packages/core/src/db/connection.ts、packages/core/src/db/adapters/sqlite.ts、packages/core/src/db/adapters/postgres.ts |
| 配置 | packages/core/src/config/config-loader.ts |
| SSE 流 | packages/server/src/adapters/web/transport.ts、packages/server/src/adapters/web/workflow-bridge.ts |
| Web UI hooks | packages/web/src/hooks/useSSE.ts、packages/web/src/lib/api.ts |
结语:一条消息的完整旅程
把全文串起来看,一条 Slack 消息的完整旅程是:入口授权 → 锁管理排队 → 编排器路由(5 个白名单命令走确定性分发,其余交给 AI)→ 工作流派发(解析隔离 → 执行器按 DAG/Loop/Step 模型运行)→ 事件发射(经 WorkflowEventBridge 转发)→ SSE 推送到 Web UI(50ms 缓冲 + 5s 重连宽限);全程的会话上下文由不可变的 Session 状态机管理,数据落库由 IDatabase 抽象统一适配 SQLite/PostgreSQL,配置由 4 层合并按优先级生效。
理解这套链路后,无论是排查"为什么消息没被 AI 处理"(看路由分叉与 prompt 组装)、"为什么工作流跑在哪个 worktree 上"(看 7 步隔离算法)、还是"为什么前端丢了实时输出"(看重连宽限与事件桥),都可以快速定位到对应的源码文件,这也是 .claude/docs/architecture-deep-dive.md 这份文档与本文存在的价值所在。
【免费下载链接】ArchonThe first open-source harness builder for AI coding. Make AI coding deterministic and repeatable.项目地址: https://gitcode.com/GitHub_Trending/archon3/Archon
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考