ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Archon 架构深度剖析:从 Slack 消息到工作流执行的全链路数据流

Archon 架构深度剖析:从 Slack 消息到工作流执行的全链路数据流 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.tsSlack event → SlackAdapter.start() 注册 app_mention message 处理器 → 授权检查isSlackUserAuthorized()auth.ts:27 → void this.messageHandler(event) —— fire-and-forgetadapter.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:392message事件只处理channel_type im的私聊并跳过 bot 自身消息以避免死循环adapter.ts:421-433。静默拒绝未授权用户授权检查失败时仅记录掩码日志如U123***不回复任何内容避免向未授权用户暴露 bot 存在adapter.ts:395-399。斜杠命令也走同一链路/archon与/archon-workflow两个 Slack 斜杠命令adapter.ts:461-468会先发送一条可见的 seed 消息作为线程根再把转换后的文本注入同一messageHandleradapter.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则负责从回复中提取这些结构化指令。源码中还有针对模型输出格式的防御逻辑例如normalizeCommandTextorchestrator-agent.ts:401-403会剥离**\/register-project ...**这类被 markdown 加粗污染的命令行保证isCommandFullyParsed能正确识别。1.4 关键决策点小结决策点行为getStreamingMode()Slack 返回batchWeb 返回streambuildFullPrompt()有 codebase 用项目级 prompt否则用全局 orchestrator promptparseOrchestratorCommands()由 AI 决定派发工作流还是纯对话回复Session resume将session.assistant_session_id传入 SDK 的options.resume确定性命令数量仅 5 个其余一切含斜杠命令均由 AI 路由2. 工作流执行/workflow run archon-fix-github-issue #422.1 派发链路用户消息以 /workflow 开头 → commandHandler.handleCommand()orchestrator-agent.ts:422 → discoverWorkflowsWithConfig() 按名称找到工作流workflows/src/loader.ts → 返回 CommandResultresult.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-forgetorchestrator.ts:336Web 与 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 → SingleStepexecuteStepInternal() → loadCommandPrompt(cwd, commandName) —— 先搜仓库再回落 bundled defaults → substituteWorkflowVariables() —— 替换 $ARGUMENTS、$ARTIFACTS_DIR 等变量 → withIdleTimeout(aiClient.sendQuery(), idleTimeout) → 流式或批量输出到平台 → ParallelBlockPromise.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 UIdag_node、workflow_step等驱动前端的工作流进度卡片实时刷新。3. 隔离解析7 步 Worktree 算法3.1 解析入口validateAndResolveIsolation()orchestrator.ts:108 → IsolationResolver.resolve(request)isolation/src/resolver.ts:100IsolationResolver按严格优先级执行 7 步探测尽量复用已有环境避免无谓地新建 worktreepackages/isolation/src/resolver.ts步骤探测内容结果1store.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 }5PR 分支收养findWorktreeByBranch(canonicalPath, prBranch)命中 →store.create({ adopted: true })→{ method: branch_adoption }6数量上限store.countActiveByCodebase()vsmaxWorktrees25触顶 →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机制非常实用如果目标分支上已经存在一个旧 worktreeArchon 不会重复创建而是直接接管它从而保留现场、节省磁盘并避免冲突。4. Session 生命周期状态机最重要的设计原则Session 迁移是不可变的——已有 session 永远不会被修改只会被停用并替换。这样每段会话的历史assistant_session_id、上下文都保留在只读记录中可随时回溯。4.1 典型迁移流程首条消息 → transitionSession(first-message) → INSERT 新 sessionparent_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) → 停用当前 sessionended_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: RecordTransitionTrigger, 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并立即创建新 sessiondeactivates仅停用当前 session由下一条消息触发新 sessionnone不执行任何操作。该 Record 类型保证了编译期穷尽性——新增触发器若未归类会直接触发 TypeScript 编译错误。这一设计在源码注释中被明确为单一事实来源single source of truth。4.3 审计追踪getSessionChain(sessionId)通过递归 CTE 沿parent_session_id链接回溯整个会话链把这次对话经历过哪些上下文切换完整还原出来——这既是调试会话上下文问题的利器也是审计为什么 AI 在某个节点丢失上下文的依据。5. 数据库层IDatabase 抽象5.1 自动检测connection.ts:30-46DATABASE_URL 已设置 → PostgresAdapterpg.Poolmax: 10 否则 → SqliteAdapterbun:sqliteWAL 模式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占位符原生可用SQLiteconvertPlaceholders()把$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 方言差异对照表FeatureSQLitePostgreSQLnow()datetime(now)NOW()jsonMerge(col, $N)json_patch(col, $N)col \|\| $N::jsonbUUIDcrypto.randomUUID()gen_random_uuid()6. 配置加载4 层合并配置不是单一文件而是按优先级从低到高合并 4 层packages/core/src/config/config-loader.tsLayer 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联动。工作流模型解析优先级当工作流执行需要解析模型时按以下顺序回落节点级modelDAG 模式per-node 声明工作流级modelYAML 顶层声明配置assistants.{provider}.modelSDK 默认值。在聊天场景下源码 orchestrator-agent.ts 的resolveChatModelRequest还引入了更高优先级的 per-user 偏好default_modeldefault_provider匹配时才生效并以tier层级别名机制做整体回落——这说明模型解析在聊天与工作流两条路径上是分别实现的工作流路径保持large层级语义。7. Web UI 数据流React → SSE → ServerWeb 端采用RESTTanStack Query v5承载静态数据 SSE 承载实时增量的双通道架构。7.1 REST 数据TanStack Query v5React 组件 → useQuery({ queryKey, queryFn }) → apiClient.listConversations() —— fetch(/api/conversations) → ServerHono 路由处理器 → DB 查询 → JSON 响应 → TanStack Query 负责缓存、轮询、失效7.2 SSE 实时流ReactuseSSE(conversationId)web/src/hooks/useSSE.ts → new EventSource(${SSE_BASE_URL}/api/stream/${conversationId}) → ServerstreamSSE(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 Patterns8.1 Lazy Logger延迟日志器每个模块都延迟创建 logger避免测试 mock 时机问题let cachedLog: ReturnTypetypeof 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-ChannelIPlatformAdapter.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.tsAI Providerpackages/providers/src/claude/index.ts、packages/providers/src/registry.ts命令处理packages/core/src/handlers/command-handler.tsSessionpackages/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.tsSSE 流packages/server/src/adapters/web/transport.ts、packages/server/src/adapters/web/workflow-bridge.tsWeb UI hookspackages/web/src/hooks/useSSE.ts、packages/web/src/lib/api.ts结语一条消息的完整旅程把全文串起来看一条 Slack 消息的完整旅程是入口授权 → 锁管理排队 → 编排器路由5 个白名单命令走确定性分发其余交给 AI→ 工作流派发解析隔离 → 执行器按 DAG/Loop/Step 模型运行→ 事件发射经 WorkflowEventBridge 转发→ SSE 推送到 Web UI50ms 缓冲 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),仅供参考
返回列表