genworker

genworker 架构说明

1. 项目概览

genworker 是一个面向组织场景的数字员工运行时。系统把来自 HTTP、IM、事件和定时任务的输入统一收敛到 WorkerRouter,再结合 Worker 主体、Duty 职责、Goal 目标、Skill 匹配、上下文管理、执行引擎、工具管线、治理与记忆系统完成一次可追踪执行。

从工程实现看,它保留多租户、多 Worker、Skill、Task、Channel 等 Agent 编排能力;从产品和领域模型看,这些能力由 Worker、Duty、Goal、Run、Process、Approval、Audit、TaskManifest、ConversationSession、InboxItem、Episode、Rule 等对象共同组成数字员工运行时。

对象模型文档中的“当前”指当前已经确定的设计目标,不要求每个对象都已经完整代码落地;未完整实现的对象会以目标边界和实现状态的方式描述。

数字员工领域对象模型的详细定义见:

当前支持四类运行入口:

技术栈

层级 技术选型
HTTP 与流式输出 FastAPI + Uvicorn + SSE
执行引擎 ReAct / Workflow / Hybrid
模型路由 LiteLLM
工具协议 MCP
持久化 Redis + MySQL + workspace 文件系统 + OpenViking 派生索引
记忆后端 OpenViking(统一检索索引)
调度 APScheduler
数据模型 @dataclass(frozen=True) 为主

补充边界:

1.1 运行时底座

运行时底座统一了后端选择、fallback 和对外状态表达:

2. 整体架构

2.1 分层架构图

flowchart LR
    subgraph Entry["入口层"]
        API["src/api\nHTTP / SSE"]
        DASH["src/services/dashboard\nWorker Dashboard 聚合"]
        CH["src/channels\nIM webhook / outbound transport"]
        EVT["src/events\nEventBus"]
        SENS["src/worker/sensing\npoll / push sensors"]
        SCHED["APScheduler\nheartbeat / goal / duty"]
        CONV["src/conversation\nthread sessions"]
        AUTO["src/autonomy\ninbox / main session / isolated run"]
    end

    subgraph Runtime["核心运行时"]
        ROUTER["src/worker/router.py\nWorkerRouter"]
        TRUST["src/common/tenant.py\nTenantLoader + TrustLevel"]
        WREG["src/worker/registry.py\nWorkerRegistry"]
        SKILL["src/skills\nSkillRegistry + SkillMatcher"]
        CTX["src/context\nContext Assembler + Compaction"]
        ENGINE["src/engine\nReact / Workflow / Hybrid"]
        TASK["src/worker/task_runner.py\nTaskRunner"]
        TOOLS["src/tools\nMCPServer + ToolPipeline + Sandbox"]
        MEMORY["src/memory\nSemantic / Episodic / Preference"]
    end

    subgraph Infra["基础设施层"]
        BOOT["src/bootstrap\nBootstrapOrchestrator"]
        LLM["src/services/llm\nLiteLLM"]
        EXT["src/services\nFeishu / WeCom / DingTalk / Slack / Email / HTTP"]
        REDIS["Redis"]
        OV["OpenViking\n派生记忆索引"]
        MYSQL["MySQL"]
        WS["workspace/\n租户与运行时文件"]
    end

    API --> DASH
    API --> ROUTER
    DASH --> WS
    DASH --> AUTO
    CH --> ROUTER
    EVT --> ROUTER
    SENS --> AUTO
    SCHED --> AUTO
    CONV --> ROUTER
    AUTO --> ROUTER

    ROUTER --> TRUST
    ROUTER --> WREG
    ROUTER --> SKILL
    ROUTER --> CTX
    ROUTER --> TASK
    CTX <--> MEMORY
    TASK --> ENGINE
    ENGINE --> TOOLS
    ENGINE --> LLM
    TOOLS --> EXT

    ROUTER <--> WS
    CONV <--> REDIS
    AUTO <--> REDIS
    MEMORY <--> OV
    TOOLS --> MYSQL
    BOOT --> API
    BOOT --> CH
    BOOT --> SENS
    BOOT --> ROUTER
    BOOT --> TOOLS

2.2 内部主要领域与关系

领域 主要模块 关系
接入层 src/api/src/channels/ 接收外部请求或消息;channels 同时统一出站 transport 适配器
会话层 src/conversation/ 维护 thread session 与会话恢复
自治运行时 src/autonomy/ 维护 main session、heartbeat inbox、isolated run 等后台认知运行时
数字员工对象层 src/worker/models.pysrc/worker/duty/src/worker/goal/src/worker/task.py 承载 Worker 主体、Duty 职责、Goal 目标、Task 执行快照
Worker 运行时 src/worker/ 负责租户加载、Worker 选择、任务级治理、Skill 路由、上下文拼装、任务提交
执行层 src/engine/src/streaming/ 按策略模式执行任务,并把内部事件编码为 SSE/AG-UI 事件流
工具层 src/tools/ 提供工具注册、发现、权限过滤、中间件、执行审计与沙箱
记忆与上下文 src/memory/src/context/ 检索长期/情景记忆,控制 token 预算,并对历史进行渐进压缩
感知与自治 src/worker/sensing/src/worker/heartbeat/src/worker/duty/src/worker/goal/ 把外部变化转为事实、职责任务、目标检查或主会话认知输入
治理与审批 src/worker/governance/src/worker/trust_gate.pysrc/tools/pipeline.pysrc/worker/lifecycle/ WorkerRouter 前置任务治理、租户信任、工具 effect 过滤、confirmation 和 suggestion approval
平台服务 src/services/ 封装 LLM、数据库、Redis、邮件和 IM 平台客户端;共享注册与生命周期模板放在 src/services/client_registry.py
启动编排 src/bootstrap/ 负责依赖排序、初始化顺序和应用生命周期
Lifecycle Intelligence src/worker/lifecycle/ 维护 task/goal/duty provenance、suggestion、feedback、detector 与 goal projector 闭环

2.2.1 执行引擎边界

当前执行引擎分为五档:

推荐选型:

2.3 脚本执行通道

执行层现在额外支持“脚本化一次工具调用”:

统一原则:

当前关键路径:

2.3 LLM 分层调度

LLM 调度拆成四层:

  1. 调用点用 src/services/llm/intent.pyLLMCallIntent 声明业务语义。
  2. src/services/llm/routing_policy.pyTableRoutingPolicyPurpose + flags 选择 tier。
  3. src/services/llm/router_adapter.py 通过 tier_aliases 把 tier 解析为 LiteLLM model group。
  4. LiteLLM Router 负责同组负载均衡、重试和 fallback。

当前 tier 包括 fast / standard / strong / reasoning,并为 tool-use 单独维护 四个 base tier。requires_tools 保留在 LLMCallIntent 中,但由 TableRoutingPolicy 内部消化;其中 fast + tools 会软升级到 standard, 避免轻量模型承担不稳定的工具调用。默认 fallback 链为:

适配器会为每次请求记录 purposetier、请求 model group 和实际响应 model, 用于成本归因、fallback 观察和延迟分析。

补充说明:

3. 主要数据流程

3.1 主数据流程图

flowchart TD
    INPUT["HTTP / IM / EventBus / APScheduler / Sensor"] --> NORMALIZE["入口标准化\nrequest / event / inbox item"]
    NORMALIZE --> SESSION["SessionManager / autonomy.MainSessionRuntime / autonomy.SessionInboxStore"]
    SESSION --> ROUTER["WorkerRouter.route_stream()"]
    ROUTER --> TENANT["加载 Tenant 与 Worker\ntrust gate / worker registry"]
    TENANT --> GOV["TaskGovernanceStage\nmanifest normalize / risk / duty binding / decision"]
    GOV --> SKILL["SkillMatcher\n使用治理后的 skill hint / preferred_skill_ids"]
    SKILL --> TOOLSCOPE["tool_scope\n按 governance posture 过滤工具"]
    TOOLSCOPE --> CONTEXT["build_managed_context()\nidentity / rules / contacts / provenance"]
    CONTEXT --> ENGINE["EngineDispatcher\nReact / Workflow / Hybrid"]
    ENGINE --> TOOLPIPE["ToolPipeline\nhooks -> governance check -> middlewares -> executor"]
    TOOLPIPE --> TOOLCALL["内置工具 / MCP 工具 / 外部服务"]
    TOOLCALL --> ENGINE
    ENGINE --> RESULT["文本结果 / 结构化事件 / 错误"]
    RESULT --> STREAM["EventAdapter / AG-UI / SSE"]
    RESULT --> PERSIST["TaskStore / SessionStore / Memory / EventBus"]

3.2 各入口如何进入统一管线

入口 首层模块 进入统一执行链路的方式
对话流 src/api/routes/chat_routes.py 先经 SessionManager,再进入 WorkerRouter.route_stream()
任务流 src/api/routes/worker_routes.py 直接进入 WorkerRouter.route_stream()
IM 双向消息 src/channels/router.py 路由为对话消息或感知事实,再进入会话层、自治运行时或 Worker 管线;当前支持 feishu / wecom / dingtalk / email / slack
事件驱动 src/events/ + src/worker/duty/ 事件匹配 Duty 后提交给 WorkerScheduler / WorkerRouter
Sensor / Heartbeat / Goal 检查 src/worker/sensing/ + src/worker/heartbeat/ + src/autonomy/ 先写入 SessionInboxStore,再由 HeartbeatRunner 决定总结、任务或 isolated run

3.3 一次执行中的关键决策点

  1. 入口先确定 tenant_idworker_idthread_idmain_session_key,并把 source_type/source_id 作为入口事实传给 Router。
  2. TenantLoaderWorkerRegistry 决定租户边界、默认 Worker 与 trust gate。
  3. TaskGovernanceStage 在 Skill 匹配前归一化 TaskManifest,执行风险分层、Duty 绑定、合同快照、治理策略快照、治理裁决和 provenance 补齐;策略快照写入 runtime/policy_snapshots/task_governance/ 并通过 policy_snapshot_refsnapshot:policy:{ref} policy ref 进入 manifest。
  4. SkillMatcher 在显式 skill_id、治理补齐的 Duty skill hint、软偏好 preferred_skill_ids 与自然语言意图之间做最终匹配。
  5. tool_scopegovernance_action + execution_posture 过滤工具;proposal_only / degraded / draft / propose 只暴露 read / compute 类工具。工具 effect 优先读取 schema 显式字段,其次读取 workspace tool_effects.yml/.yaml sidecar,再读取 Worker / Tenant tool_effect_overrides,未知动态工具默认高风险。
  6. build_managed_context() 组装身份、约束、规则、联系人、情景记忆、历史消息和治理后的 provenance。
  7. EngineDispatcher 根据 Skill 策略选择 ReAct、Workflow 或 Hybrid。
  8. 工具调用统一经过 ToolPipeline,再做一次 governance effect 校验,防止运行期动态工具绕过 Router 过滤。
  9. 如果 TaskManifest.pre_script 非空,TaskRunner 会先执行脚本并把 stdout 注入任务文本;生产路径中首次 manifest 创建和治理应发生在 Router。
  10. execute_code / ScriptTool 而言,子进程内的工具调用会通过 RPC bridge 回到同一条 ToolPipeline
  11. 结果同步写回任务、会话、记忆和事件流。
  12. post-run lifecycle hook 会补充 task 与 goal/duty 的显式关联,并在满足条件时生成 pending suggestion。

spawn_taskdelegate_to_workerspawn_subagents 还会在工具 handler 内读取当前 ExecutionScope.governance 做二次校验:父任务没有 delegate effect 权限时,即使工具被直接调用也会拒绝派生执行。

新增工具必须在 schema、workspace tool_effects.yml/.yaml 或 Worker / Tenant tool_effect_overrides 中声明 effect / risk_level / requires_approval;未声明的运行期工具按 high-risk external commit 处理,不会暴露给 proposal_only 或降级姿态。

3.4 Lifecycle 联动

新增的 lifecycle intelligence layer 现在挂在现有执行链路上,而不是单独开一条执行路径:

3.4 对话/任务请求时序图

sequenceDiagram
    participant Client as Client
    participant Route as chat_routes / worker_routes
    participant Ingress as prepare_service_ingress
    participant Session as SessionManager
    participant Router as WorkerRouter
    participant Context as build_managed_context
    participant Engine as EngineDispatcher
    participant Tools as ToolPipeline
    participant Stream as EventAdapter / SSE

    Client->>Route: POST /api/v1/chat/stream or /worker/task/stream
    Route->>Router: resolve_entry(task, tenant_id, worker_id)
    alt service 模式
        Route->>Ingress: prepare_service_ingress(...)
        Ingress->>Session: get_or_create(service thread)
        Ingress-->>Route: session_ttl / metadata / task_context
    end
    opt 对话模式
        Route->>Session: get_or_create(thread session)
        Route->>Session: append user message
    end
    Route->>Router: route_stream(...)
    Router->>Context: assemble identity / rules / contacts / episodic / history
    Router->>Engine: dispatch by strategy mode
    loop ReAct / Workflow / Hybrid
        Engine->>Tools: execute tool calls when needed
        Tools-->>Engine: tool results
    end
    Engine-->>Router: stream events / text / errors
    Router-->>Route: StreamEvent
    Route->>Stream: format as AG-UI / legacy SSE
    Stream-->>Client: text/event-stream
    opt 对话模式
        Route->>Session: save updated session
    end

说明:

3.5 感知与 Heartbeat 时序图

sequenceDiagram
    participant Sensor as Sensor / Channel Router / Event Source
    participant Inbox as SessionInboxStore
    participant Bus as EventBus
    participant HB as HeartbeatRunner
    participant Main as MainSessionRuntime
    participant Strategy as HeartbeatStrategy
    participant Scheduler as WorkerScheduler / IsolatedRunManager
    participant Router as WorkerRouter

    Sensor->>Inbox: write(InboxItem)
    Inbox->>Bus: publish(inbox.item_written)
    Note over HB: APScheduler 周期触发 run_once()
    HB->>Inbox: fetch_pending(tenant_id, worker_id)
    Inbox-->>HB: PENDING -> PROCESSING items
    HB->>HB: dedupe via AttentionLedger
    HB->>Strategy: decide_action(item)
    alt summary
        HB->>Main: build_task_context()
        HB->>Router: route_stream(summary prompt)
        Router-->>HB: assistant summary
        HB->>Main: append heartbeat message
    else task
        HB->>Scheduler: submit_task(manifest, preferred_skill_ids)
    else isolated
        HB->>Scheduler: create isolated run(main_session_key)
    end
    HB->>Inbox: mark_consumed(...)
    HB->>Main: update_heartbeat_state(...)

说明:

4. 核心模块说明

4.1 API 与流式输出

模块 职责
src/api/app.py FastAPI 应用工厂、生命周期入口、路由装配
src/api/routes/chat_routes.py 对话流式入口与任务查询接口
src/api/routes/worker_routes.py Worker 任务流式入口、运维概览、Worker 局部重载
src/api/routes/worker_dashboard_routes.py Worker 概览、任务、日历、Goal、活动、待办和渠道只读查询
src/api/routes/channel_routes.py IM webhook 接入与渠道状态查询
src/api/routes/employee_routes.py 数字员工对外 Facade(src/services/employee/),按 employee_id 暴露能力/运行查询
src/api/routes/openai_compat_routes.py OpenAI 兼容适配层(src/services/openai_compat/),把外部 chat/completions 请求归一化后进入主链路
src/streaming/ 内部 StreamEvent 到 SSE/AG-UI 协议的编码适配

Dashboard 查询路径位于 src/services/dashboard/,按 profile / task / calendar / goal / activity / pending-actions / channel / runs / capabilities 拆分领域 service; 聚合器只做并发调用和单 block 失败降级,不反向进入各领域 service。runs block 经 RunRecordStore 读取 RunRecord.deliverables(设计 §2.3,旧 result_summary[:500] 截断 不再是唯一交付物出口),capabilities block 复用 RuntimeCapabilityRegistry 而不直接读 散落的 app.state 组件字段(设计 §5.2)。

Dashboard 当前是 Worker 运行状态聚合视图,按 profile、task、calendar、goal、activity、pending-actions、channel、runs 和 capabilities 组织数据。

4.1.1 对外入口与审计闭环

外部入口(Employee Facade / OpenAI 兼容适配器)不绕过统一执行链路:它们在系统边界做协议归一化与鉴权后,仍进入 WorkerRouter.route_stream() 主链路。执行收尾时 TaskRunner 通过 RunRecordBuilder(采集埋点 run_collectors.py)把工具调用、数据访问、治理裁决与交付物落盘为 RunRecord,由 RunRecordStore 持久化到 .../workers/{wid}/runs/。对外的运行/交付物查询(Employee API)与 Dashboard runs block 再从同一份 RunRecord 读取,形成「外部入口 → 主链路 → RunRecord 落盘 → 审计/交付物回读」的闭环;能力查询则统一经 RuntimeCapabilityRegistry 投影,与 /api/v1/runtime/capabilities 共用同一状态源。

4.2 Worker 运行时

模块 职责
src/worker/models.py Worker、ServiceConfig、HeartbeatConfig、ToolPolicy 等核心定义
src/worker/parser.py 解析 PERSONA.md frontmatter 和正文
src/worker/registry.py Worker 注册、默认 Worker 与关键词匹配
src/worker/loader.py Worker 磁盘装载、技能合并和标准运行时目录补齐
src/worker/router.py 统一编排入口,负责 worker/skill 决策并驱动执行
src/worker/tool_scope.py 组装 per-run tool bundle、task tools、session search 与 execution scope
src/worker/runtime_context.py 装载 learned rules、episodic/profile/preference/decision 上下文并生成 WorkerContext
src/worker/task_runner.py 负责一次运行的执行包装、回调与学习反馈
src/worker/scheduler.py 统一接收 duty、goal、heartbeat、isolated run 产生的任务
src/worker/scheduler_effects.py 处理调度完成事件、isolated run 通知与 dead-letter 持久化
src/worker/trust_gate.py 基于租户信任等级决定 bash、remote MCP、episodic write 等能力边界

说明:

4.3 执行引擎

模块 职责
src/engine/router/engine_dispatcher.py 根据 Skill 策略模式派发引擎
src/engine/react/agent.py ReAct 自主推理与工具调用循环
src/engine/workflow/engine.py 确定性工作流执行
src/engine/hybrid/engine.py 组合自主步骤与确定性步骤
src/engine/tools/subagent_tool.py Worker 内部 SubAgent 并行编排工具

4.4 工具框架

模块 职责
src/tools/mcp/server.py 工具注册中心
src/tools/pipeline.py 工具执行管线
src/tools/sandbox.py 工具权限过滤
src/tools/middlewares/ 权限、schema 校验、超时、审计、脱敏
src/tools/builtin/ bash_execute、文件读写、搜索、网页抓取、agent 委托等内置工具

4.5 上下文与记忆

模块 职责
src/context/assembler.py 组装最终上下文片段
src/context/budget_allocator.py 分配 system、history、rules、memory 等 token 配额
src/context/compaction/ 负责 tool trim、history prune、history summarize、reactive recovery
src/memory/orchestrator.py 统一语义记忆、情景记忆、偏好提取与故障隔离
src/memory/episodic/ 情景记忆存储、检索、衰减与关联反馈
src/memory/preferences/ 偏好与长期决策抽取

说明:

4.6 会话、感知与自治

模块 职责
src/conversation/session_manager.py thread session 与 main session 生命周期管理
src/autonomy/inbox.py 感知事实的三态 Inbox 存储与原子消费
src/autonomy/main_session.py 主会话状态、heartbeat 元数据和 isolated run 汇总
src/autonomy/isolated_run.py 隔离执行任务创建与回流
src/worker/heartbeat/runner.py 消费 Inbox,决定总结、任务或 isolated run
src/worker/sensing/registry.py 管理 Email、Webhook、Git、Workspace File 等 sensor
src/worker/duty/trigger_manager.py EventBus 事件到 Duty 的匹配与调度
src/worker/goal/progress_checker.py Goal 偏差检测与进展评估

说明:

4.7 Channel Transport

模块 职责
src/channels/outbound.py 统一 Email/Feishu/WeCom/DingTalk 出站适配器、重试与多通道回退
src/channels/outbound_types.py ChannelMessageSenderScopeRetryConfig 等共享出站模型
src/channels/bindings.py worker-scoped channel binding 构造与凭据过滤 helper
src/channels/router.py IM 入站总入口,负责命令判定、对话转发、sensor 路由和完成通知订阅
src/worker/integrations/worker_scoped_channel_gateway.py tenant_id + worker_id + channel_type 解析 worker-scoped 出站 transport

说明:

4.8 Data Access

统一微盘挂载把飞书云空间、企微微盘和钉盘收敛为 DriveClient 协议,由 MountManager 统一路由到 worker-scoped 平台 client。PERSONA.mdmounts[] 声明挂载,src/bootstrap/mount_init.py 在启动期为 worker 构造 MountManagersrc/runtime/drive_runtime/ 订阅 drive sensor 事件并把远端文件镜像到 data/scratch/drives/{mount_id}/

模块 职责
src/services/_drive/protocol.py DriveClient / DriveFileMetadata 统一协议
src/worker/data_access/mount_registry.py drive sensor 只读查询 MountConfig 的协议边界
src/worker/data_access/mount_manager.py mounts/{mount_id}/... 虚拟路径读写、token refresh 与 cache invalidation
src/worker/sensing/sensors/drive_sensor_base.py 三平台 drive sensor 的 snapshot diff、过滤、删除与重命名检测
src/runtime/drive_runtime/mirror.py drive runtime 门面:事件订阅、手动 reconcile、retry drain、sidecar rebuild
src/runtime/drive_runtime/mirror_executor.py event 与 reconcile 共用的单文件 mirror/export/delete 执行体
src/runtime/drive_runtime/sync_coordinator.py sync state、retry queue、failed queue 与领域事件发布的单一入口
src/runtime/drive_runtime/reconcile.py 启动期 full reconcile 与周期性 audit 的 remote/local diff
src/runtime/drive_runtime/text_extractors.py .docx/.xlsx 本地文本 sidecar 生成与读取
src/runtime/drive_runtime/index.py mirror 完成后优先读取 sidecar 写 episodic memory,并在可用时写 OpenViking semantic index
src/api/drive_config_platforms.py Drive 配置平台规格注册表,收敛 wecom / feishu / dingtalk 的 source、credentials、sensor type 与 client builder 差异
src/api/drive_config_loader.py 读写 worker PERSONA.md 中的 mounts[]sensor_configs[],组合保存时提交 PERSONA.md + CHANNEL_CREDENTIALS.json 两文件快照
src/api/routes/drive_config_routes.py GET/PUT/DELETE /api/v1/workers/{worker_id}/drive-config... 配置入口;路由层只做 HTTP 转换、鉴权、异常映射
src/api/routes/drive_sync_settings_routes.py GET/PUT /api/v1/drives/sync-settings workspace 级 sync settings
src/api/_channel_credentials_writer.py IM 与 drive 共用的 worker CHANNEL_CREDENTIALS.json 原子写入口
src/api/webhook_ingress/ /api/v1/webhooks/{provider} 统一 provider 验签与 kind dispatch
src/tools/builtin/drive_search_tool.py 对镜像内容执行 grep / semantic / hybrid 搜索

Drive sensor 的 filter 只保存 mount_id 和 sensor 自有字段(例如 extensions)。平台 source 的唯一来源是 mounts[].source:企微为 space_id / folder_id / space_type,飞书为 folder_token / space_type,钉钉为 union_id / path / space_type。运行时由 sensor 通过注入的 MountRegistry.get_mount(mount_id) 反查,避免手工编辑时出现 source 与 filter 漂移。

Drive mirror 的持久状态按 tenant + worker + mount 隔离在 data/scratch/drives/{mount_id}/.sync_state.json.sync_queue.json.failed.jsondrive.mirrored 是唯一稳定成功事件;在线文档以 .docx/.xlsx 原生文件落盘,同时写 <local_file>.text.json sidecar 供 DriveIndexServicedrive_search_tool 本地检索,搜索路径不调用平台 API。

Drive config API 以配置聚合为写入口。组合保存 mount 与 credentials 时会先生成 PERSONA.mdCHANNEL_CREDENTIALS.json 的完整新快照,再提交文件;提交失败不会触发 worker reload。reload_worker_runtime_state()PERSONA.md 变更路径中先调用 refresh_mount_manager() 重建该 worker 的 MountManager,再刷新 sensor registry,确保新 sensor 注入的是最新 mount source。单独保存 CHANNEL_CREDENTIALS.json 时只失效 credential loader 与 platform client factory,不重建 mount manager。

4.9 启动编排

src/bootstrap/ 使用 BootstrapOrchestrator 按依赖顺序初始化系统。当前主链路包含:

logging
  -> events
  -> llm
  -> memory
  -> mcp
  -> tool_discovery
  -> skills
  -> workers
  -> api_wiring
  -> scheduler
  -> conversation
  -> platforms
  -> mounts
  -> contacts
  -> contact_sync
  -> channels
  -> integrations
  -> sensors
  -> drive_runtime
  -> heartbeat

其中 src/runtime/ 继续吸收原本堆积在 composition root 的应用服务逻辑:

这样 src/bootstrap/__init__.py 主要负责初始化顺序与模块聚合;src/bootstrap/api_wiring_init.py 承载 initializer 壳,src/bootstrap/compat.py 承载历史兼容导出,根模块仅通过懒加载转发旧入口,不再在根模块堆放大块业务实现。

4.10 平台服务边界

模块 职责
src/services/client_registry.py 提供 keyed singleton registry、显式初始化和按数据库名关闭等共享生命周期模板
src/services/mysql/client.py 实现 MySQL 连接池、查询/事务协议,并复用共享 registry 处理多库实例管理
src/services/redis/client.py 实现 Redis 连接池与 KV/集合操作,并复用共享 registry 处理默认实例生命周期

说明:

4.11 兼容层约束

当前仓库仍保留少量历史导入路径,但这些模块都属于 external-only compatibility layer:

约束:

更准确地说,运行时启动顺序是依赖拓扑,而不是严格单链。核心依赖关系如下:

flowchart LR
    LOG["logging"] --> EVT["events"]
    LOG --> LLM["llm"]
    LOG --> SK["skills"]
    LLM --> MEM["memory"]
    MEM --> MCP["mcp"]
    MCP --> TD["tool_discovery"]
    SK --> WK["workers"]
    WK --> APIW["api_wiring"]
    APIW --> SCH["scheduler"]
    EVT --> SCH
    APIW --> CONV["conversation"]
    EVT --> CONV
    SCH --> CONV
    WK --> PLAT["platforms"]
    WK --> CT["contacts"]
    WK --> INT["integrations"]
    APIW --> INT
    SCH --> INT
    WK --> CH["channels"]
    APIW --> CH
    CONV --> CH
    PLAT --> CH
    EVT --> CH
    CT --> CH
    SCH --> SENS["sensors"]
    EVT --> SENS
    PLAT --> SENS
    INT --> SENS
    SCH --> HB["heartbeat"]
    CONV --> HB
    INT --> HB

5. 运行时与数字员工对象结构

本章同时描述两层对象:

详细字段和关系见 docs/DIGITAL_EMPLOYEE_OBJECT_MODEL.md。本章只保留架构级关系,避免把架构文档膨胀成完整对象说明。

5.1 数字员工对象关系图

erDiagram
    TENANT ||--o{ WORKER : contains
    WORKER ||--o{ SKILL : loads
    WORKER ||--o{ DUTY : owns
    WORKER ||--o{ GOAL : pursues
    WORKER ||--o{ TASK_MANIFEST : executes
    WORKER ||--o{ CONVERSATION_SESSION : serves
    WORKER ||--o{ INBOX_ITEM : receives
    WORKER ||--o{ EPISODE : accumulates
    WORKER ||--o{ RULE : applies
    DUTY ||--o{ DUTY_TRIGGER : triggered_by
    DUTY ||--|| EXECUTION_POLICY : uses
    DUTY }o--o| ESCALATION_POLICY : escalates_with
    DUTY ||--o{ DUTY_EXECUTION_RECORD : logs
    GOAL ||--o{ MILESTONE : contains
    MILESTONE ||--o{ GOAL_TASK : contains
    TASK_MANIFEST ||--|| TASK_PROVENANCE : carries
    TASK_PROVENANCE }o--o| GOAL : links_goal
    TASK_PROVENANCE }o--o| DUTY : links_duty
    CONVERSATION_SESSION ||--o{ CHAT_MESSAGE : contains
    TASK_MANIFEST }o--|| CONVERSATION_SESSION : references
    INBOX_ITEM }o--|| CONVERSATION_SESSION : targets

说明:

5.2 核心领域与运行时对象

对象 模块 作用
Tenant src/common/tenant.py 多租户隔离边界,提供 trust level、默认 Worker、租户级工具策略
Worker src/worker/models.py 数字员工岗位主体;承载身份、模式、工具策略、service/heartbeat/sensor/channel 配置
Duty src/worker/duty/ 职责或流程步骤执行单元;负责触发、执行策略、职责合同、流程步骤映射和执行记录
DutyDefaults src/worker/duty/ 目标模型 职责包默认规则;第一阶段 canonical 位置为 duties/_defaults.md
DutyContract src/worker/duty/ 目标模型 单个 Duty 的职责/步骤合同,包含输入输出、SLA、决策边界、升级、审批和审计要求
Goal src/worker/goal/ 阶段目标和进展管理对象;包含 milestone、goal task、优先级、状态,以及可选的 primary_duty_id / duty_refs 和 task 级 duty_id / process_id / step_id 映射
Run / TaskManifest src/worker/task.py 一次实际执行记录或生命周期快照;第一阶段 RunRecord v0 = TaskManifest + run_id + provenance + lifecycle event
TaskProvenance src/worker/task.py Task 与 goal、duty、trigger、process step、contract version、governance decision、approval、audit refs 的关联元数据
TaskGovernanceStage src/worker/governance/task_governance_stage.py WorkerRouter 的前置治理阶段,统一归一化 manifest、风险姿态、Duty 绑定、合同快照和执行动作
GovernanceDecision src/worker/governance/ 任务级治理决策,由 GovernanceDecisionResolver 统一生成,解释 auto / gated / blocked 及策略来源;proposal_only / degraded / gated_execute 由 action、posture 和 provenance 字段区分
ApprovalResolver src/worker/governance/approval_resolver.py confirmation 恢复执行时读取 approval_ref,避免 approved manifest 重复进入 gated confirmation
ContractSnapshot src/worker/duty/contract_snapshot.py 目标模型 执行时 Duty 合同快照,用于历史回放和审计
Process 目标模型 企业业务流程及流程步骤,用于把 Worker/Duty 映射到真实业务流程;当前尚未作为一等对象完整落地
Approval src/worker/lifecycle/src/autonomy/inbox.py 人工确认、suggestion approval、LangGraph interrupt 等审批请求和审批决定;统一模型仍在演进
Audit src/tools/pipeline.pysrc/worker/lifecycle/ 工具调用、治理决策、审批、执行结果和输出交付的审计记录;统一 AuditRecord 仍在演进
ConversationSession src/conversation/models.py thread/main session 的统一持久化载体
InboxItem src/autonomy/inbox.py 感知层写入的结构化事实
Event src/events/models.py 进程内事件总线的统一事件格式
ContextWindowConfig src/context/models.py 上下文预算、压缩阈值和各 segment 配额
Episode src/memory/episodic/models.py 情景记忆条目
Rule src/worker/rules/models.py directive / learned 规则实体

5.3 Workspace 布局

workspace/
├── system/
│   └── skills/
│       └── {skill_name}/SKILL.md
└── tenants/
    └── {tenant_id}/
        ├── TENANT.json
        ├── sessions/
        │   └── {session_id}.json
        ├── skills/
        │   └── {skill_name}/SKILL.md
        └── workers/
            └── {worker_id}/
                ├── PERSONA.md
                ├── CHANNEL_CREDENTIALS.json        # 可选
                ├── runtime/
                │   ├── inbox.json
                │   └── heartbeat_meta.json
                ├── skills/
                │   └── {skill_name}/SKILL.md
                ├── contacts/
                │   ├── configured/
                │   ├── discovered/
                │   └── index.jsonl
                ├── duties/
                ├── goals/
                ├── rules/
                │   ├── directives/
                │   └── learned/
                ├── memory/
                │   └── episodes/
                ├── tasks/
                │   └── active/
                ├── sensor_snapshots/
                └── archive/

说明:

5.4 Workspace 配置文件参数说明

TENANT.json

字段 类型 作用
tenant_id string 租户唯一标识,目录名与文件内容必须一致
name string 租户展示名
trust_level 0-3 租户信任等级,控制 bash、remote MCP、规则写入、episodic write 等能力边界
tool_policy.denied_tools string[] 租户级禁用工具列表,作为全租户安全覆盖层
mcp_remote_allowed boolean 是否允许远程 MCP 发现或访问
default_worker string 未显式指定 worker_id 时的默认 Worker
credentials object 租户级凭据扩展位,当前代码可读取该字段,但主要平台凭据已下沉到 Worker 目录

PERSONA.md frontmatter

PERSONA.md 由 YAML frontmatter 和 Markdown 正文组成。frontmatter 负责结构化配置,正文作为长期系统指令注入 Prompt。

顶层字段
字段 类型 作用
identity object Worker 身份、角色与人格描述
mode personal \| team_member \| service 决定上下文注入方式、会话策略与服务边界
tool_policy object Worker 级工具黑白名单
skills_dir string Worker 私有 Skill 目录,默认 skills/
default_skill string Worker 默认 Skill
constraints string[] 执行约束,注入系统 Prompt
sensor_configs object[] 传感器配置列表,决定轮询/推送事实如何进入 Inbox
mounts object[] 外部微盘挂载配置,支持 feishuwecomdingtalk
channels object[] IM 渠道绑定配置
contacts object[] 预置联系人档案
contact_settings object 联系人目录的存储与上下文注入策略
service object mode=service 时的会话和升级策略
heartbeat object Worker 级 heartbeat 判定覆盖项

兼容别名:

数字员工对象模型补充:

identity
字段 类型 作用
identity.name string Worker 展示名
identity.worker_id string Worker 唯一标识,目录名通常与其对应
identity.version string Worker 定义版本号
identity.role string 角色描述
identity.department string 所属部门
identity.reports_to string 汇报对象
identity.background string 背景知识与职责边界
identity.personality.traits string[] 人格特征
identity.personality.communication_style string 输出风格
identity.personality.decision_style string 决策偏好
identity.principles string[] 高优先级原则
tool_policy
字段 类型 作用
tool_policy.mode blacklist \| whitelist 黑名单或白名单模式
tool_policy.denied_tools string[] 禁用工具列表
tool_policy.allowed_tools string[] 白名单模式下允许的工具列表
sensor_configs[]
字段 类型 作用
source_type string 传感器类型,如 emailwebhookworkspace_filegit
poll_interval string 轮询间隔,如 5m30m
delivery_mode string 事实投递模式扩展位
filter object 传感器过滤条件
auto_create_goal boolean 是否自动由事实生成 Goal
require_approval boolean 由事实生成 Goal 或任务时是否需要人工批准
cognition_route_override string 覆盖默认认知路由
routing_rules[] object[] 按字段、模式和匹配方式将事实分类到不同 route
fallback_route string 未命中规则时的默认 route
mounts[]
字段 类型 作用
mount_id string Worker 内唯一挂载 ID,同时决定默认 mount_path=mounts/{mount_id}/
type feishu \| wecom \| dingtalk 平台类型
source object 平台源信息:飞书配置页首期支持 folder_token,企微 space_id/folder_id,钉钉 union_id/path
mount_path string 虚拟路径;解析时会归一为 mounts/{mount_id}/
permissions string[] 默认 ["read"];包含 write 时允许通过 mount 路径上传
sync_strategy string on_demandproactive,当前作为策略标记
cache_ttl number MountManager 读缓存 TTL,单位秒
source.space_type string 平台空间类型,按平台透传给 client
channels[]
字段 类型 作用
type string 渠道类型,如 feishuwecomdingtalkemail
connection_mode string 连接方式,如 webhook、polling
chat_ids string[] 绑定的会话或群组 ID
reply_mode string 回复策略,如完整回复、流式回复、卡片更新
features object 渠道特性开关,例如群聊监控
contacts[]
字段 类型 作用
person_id string 联系人唯一标识
name string 联系人主名称
role string 角色
organization string 所属组织
notes string 备注
confidence number 档案置信度
identities[] object[] 跨渠道身份,如邮箱、IM handle
social_circles string[] 社交圈分组
hierarchy_level string 组织层级标记
aliases string[] 别名
tags string[] 自定义标签
service_count number 服务模式下的互动计数
common_topics string[] 常见主题
contact_settings
字段 类型 作用
workspace_root string 联系人存储根目录
discovered_dir string 自动发现联系人目录名
configured_dir string 手工配置联系人目录名
index_file string 联系人索引文件名
context_limit number 注入 Prompt 的联系人数量上限
sync.enabled boolean 是否启用 Worker 通讯录同步
sync.interval_minutes number APScheduler 定时同步间隔
sync.delete_policy never \| mark_stale \| delete 平台缺失联系人处理策略,默认 mark_stale
sync.sources[] object[] 通讯录来源列表,支持 feishudingtalkwecomslack
sync.sources[].platform string 平台类型
sync.sources[].root_department_id string 飞书/钉钉/企微同步根部门
sync.sources[].include_child_departments boolean 是否递归子部门
sync.sources[].include_bots boolean Slack 等平台是否同步 bot
sync.sources[].include_deleted boolean 是否同步禁用/删除用户
sync.sources[].page_size number 平台分页大小
sync.sources[].max_pages number 单次同步最大页数

通讯录同步使用当前 Worker 的 CHANNEL_CREDENTIALS.json,同步结果落在 workspace/tenants/{tenant_id}/workers/{worker_id}/contacts/。同步写入通过 ContactRegistry.upsert_directory_member() 进入联系人聚合根,同一 Worker 的 contacts 目录写入使用 registry 级写锁,避免多平台并发重写 index.jsonl。configured 联系人被同步命中时保留 contacts/configured/ 原文件位置,不生成 discovered 副本。

手动同步 API:

POST /api/v1/workers/{worker_id}/contacts/sync?tenant_id=demo
GET /api/v1/workers/{worker_id}/contacts/sync/status?tenant_id=demo

POST body 示例:

{"platforms":["slack"],"dry_run":false}

运行态 debug 中会出现 contact_sync component,用于展示同步服务是否启用、最近错误和运行中任务数量。

service
字段 类型 作用
knowledge_sources[] object[] service 模式专用知识源列表
session_ttl number service 会话 TTL,单位秒
max_concurrent_sessions number 最大并发会话数
anonymous_allowed boolean 是否允许匿名访问
escalation.enabled boolean 是否允许升级到人工或其他 Worker
escalation.target_worker string 升级目标 Worker
escalation.triggers string[] 触发升级的条件描述
heartbeat
字段 类型 作用
goal_task_actions string[] 哪些 goal action 应直接转为标准 task
goal_isolated_actions string[] 哪些 goal action 应转为 isolated run
goal_isolated_deviation_threshold number goal 偏差阈值,超过后更倾向 isolated run

duties/DUTY.md

duties/ 承载 Worker 的职责定义。当前 Duty 已用于事件触发、周期调度、人工触发和职责执行记录。

设计目标上,duties/ 还承载职责合同与执行 provenance 的源头:

当前 Duty 模型字段:

字段/区域 类型 作用
duty_id string 职责唯一标识
title string 职责标题
status active \| closed \| deprecated 职责状态
triggers[] object[] schedule、event、condition、collaboration、manual 等触发条件
execution_policy object 默认执行深度和按 trigger 覆盖的执行深度
action string Markdown 形式的执行描述
quality_criteria[] string[] 职责完成质量标准
skill_id / skill_hint / preferred_skill_ids string / string[] 技能绑定或软偏好
pre_script object 执行前脚本配置
escalation object 升级条件和目标
execution_log_retention string 职责执行日志保留周期

目标 Duty 合同字段:

字段/区域 类型 作用
duty_defaults object 当前 Duty 的局部默认覆盖;共享默认规则使用 duties/_defaults.md
contract.version string 职责合同版本,执行时写入 TaskProvenance.contract_version
contract.input_schema object 本职责需要的输入字段、附件、上下文和校验要求
contract.output_schema object 本职责承诺交付的结构化结果、文本结论、事件或文件
contract.decision_scope string[] 本职责可自动判断、只能建议、必须确认或禁止的事项
contract.sla object 响应时限、完成时限、重试和超时处理
contract.escalation object 升级条件和升级目标
contract.approval object 审批触发条件、审批角色和恢复策略
contract.audit object 执行时必须记录的证据和审计要求
process_step.process_id string 业务流程 ID
process_step.step_id string 业务流程步骤 ID
process_step.step_name string 业务流程步骤展示名

执行时,DutyExecutor 应把 duty_idtrigger_idprocess_idstep_idcontract_versioncontract_snapshot_ref 和任务级 governance_decision 写入 TaskManifest.provenancecontract_snapshot_ref 指向 runtime/contract_snapshots/duties/{duty_id}/{contract_version}-{sha256_12}.json,用于历史合同回放。

治理决策由目标模型 GovernanceDecisionResolver 统一生成,合并 Tenant / Worker 上界、duties/_defaults.md、单 Duty 默认覆盖、DutyContract 和现有 resolve_gate_level() 的判断。审批结果和工具审计引用分别写入 approval_refaudit_refs,且应通过 Task lifecycle aggregate 的受控方法写回。

CHANNEL_CREDENTIALS.json

该文件位于 workspace/tenants/{tenant_id}/workers/{worker_id}/CHANNEL_CREDENTIALS.json,用于 Worker 级平台凭据装载。

feishu
字段 类型 作用
app_id string 飞书应用 ID
app_secret string 飞书应用密钥
wecom
字段 类型 作用
corpid string 企业微信 CorpID
corpsecret string 企业微信应用密钥
agent_id string 企业微信 Agent ID
dingtalk
字段 类型 作用
app_key string 钉钉应用 Key
app_secret string 钉钉应用 Secret
robot_code string 钉钉机器人编码
email
字段 类型 作用
worker_address string Worker 发信地址
worker_username string Worker 邮箱用户名
worker_password string Worker 邮箱密码、授权码或应用专用密码
worker_imap_host string Worker IMAP 主机
worker_imap_port number Worker IMAP 端口
worker_smtp_host string Worker SMTP 主机
worker_smtp_port number Worker SMTP 端口
owner_address string 所属人邮箱地址
owner_username string 所属人邮箱用户名
owner_password string 所属人邮箱密码、授权码或应用专用密码
owner_imap_host string 所属人 IMAP 主机
owner_imap_port number 所属人 IMAP 端口
owner_smtp_host string 所属人 SMTP 主机
owner_smtp_port number 所属人 SMTP 端口

邮箱集成使用 IMAP 收信、SMTP 发信,不使用 POP3。邮箱通道的 chat_ids 表示路由收件地址;单邮箱 worker 可留空,保存 IM 配置时会默认回填 worker_address。自定义邮箱目录通过通道 features.folders 显式配置;运行时工具可先调用 email_list_folders 获取真实目录名,再调用 email_search(folder=...)

5.5 持久化位置与对应对象

路径 对应对象 说明
workspace/tenants/{tenant_id}/TENANT.json Tenant 租户静态配置
workspace/tenants/{tenant_id}/sessions/{session_id}.json ConversationSession thread/main session 文件存储
workspace/tenants/{tenant_id}/workers/{worker_id}/PERSONA.md Worker Worker 定义与长期系统指令
workspace/tenants/{tenant_id}/workers/{worker_id}/duties/ Duty Worker 职责定义和触发配置
workspace/tenants/{tenant_id}/workers/{worker_id}/goals/ Goal 阶段目标、里程碑、进展和关联职责
workspace/tenants/{tenant_id}/workers/{worker_id}/tasks/active/{task_id}.json TaskManifest 任务生命周期快照
workspace/tenants/{tenant_id}/workers/{worker_id}/memory/episodes/*.md Episode 情景记忆条目
workspace/tenants/{tenant_id}/workers/{worker_id}/rules/directives/*.md Rule 指令性规则
workspace/tenants/{tenant_id}/workers/{worker_id}/rules/learned/*.md Rule 学习得到的规则
workspace/tenants/{tenant_id}/workers/{worker_id}/runtime/inbox.json InboxItem Sensor/事件写入的事实
workspace/tenants/{tenant_id}/workers/{worker_id}/runtime/heartbeat_meta.json heartbeat meta 主会话 heartbeat 游标、open concerns、task refs
workspace/tenants/{tenant_id}/workers/{worker_id}/sensor_snapshots/*.json Sensor 快照 传感器状态与增量游标

5.6 环境配置加载

系统级环境变量由 src/common/settings.py 装载,优先级如下:

  1. configs/config.env
  2. 解析后的环境文件(本地默认是 configs/config_local.env
  3. configs/config_local.env(当第 2 步不是它时,作为最终本地覆盖)
  4. 进程环境变量

入口 start.py 会先读取 ENVIRONMENTENV,再调用 load_layered_env()。当前仓库内的环境配置文件包括:

本地启动默认读取 configs/config_local.env,该文件不提交到 git,仓库提供 configs/config_local.env.example 作为模板。系统级配置负责提供 HTTP 服务、日志、MCP、Redis、MySQL、embedding、OpenViking、heartbeat、自动热刷新等参数;workspace/ 中的租户和 Worker 配置则负责业务运行面的差异化配置。

LLM 路由和 provider 配置不再放在 config.env 体系里。本地开发从 configs/litellm_local.json 读取,该文件不提交到 git,仓库提供 configs/litellm_local.json.example 作为模板;非本地环境通过启动期注入 provider 提供 LiteLLM 配置。当前 src/services/llm/config_source.py 负责区分本地文件和 非本地注入来源;后续如果改为 Nacos,只需要替换这一层 provider。

6. 架构特征

模式 应用位置 作用
Registry SkillRegistryWorkerRegistryMCPServer 支持运行时发现与组合
Pipeline ToolPipelineBootstrapOrchestrator 将多阶段处理拆分为稳定步骤
Strategy EngineDispatcherHeartbeatStrategy 在不同模式下替换执行策略
Facade WorkerRouterMemoryOrchestrator 向上提供简化入口
Protocol / 接口隔离 LLMClientToolExecutorEventBusProtocol 降低具体实现耦合
不可变数据 @dataclass(frozen=True) 避免运行时隐式副作用
渐进压缩 src/context/compaction/ 在 token 压力下保留高价值上下文
多租户安全边界 Tenant + TrustGate + ToolPolicy 统一控制高风险能力暴露

7. 测试结构

目录 用途
tests/unit/ 单模块行为验证
tests/integration/ 跨模块协同验证
tests/e2e/ 完整请求链路验证
tests/structural/ AST 结构断言,防止外层重新直接写领域状态

WorkerRouter 任务级治理相关改动至少运行:

pytest tests/unit/test_task_governance_stage.py tests/unit/test_duty_binding_resolver.py -q
pytest tests/integration/test_worker_router_task_governance.py -q
pytest tests/structural/test_worker_router_governance_boundaries.py -q

8. 领域自治调用链

Task、Goal、Inbox、Suggestion、Preference、Episodic、RuntimeState、Session 的状态变更统一经聚合根完成。外层 handler 只做 directory.find(...).action(...),不再串联 replace(...)store.save(...) 和事件发布。

Bootstrap 装配顺序中,LifecycleInitializer(95) 创建 lifecycle store 与 Suggestion directory;GoalInitializer(96) 创建 Goal directory/effects;ApiWiringInitializer(100) 创建 TaskStore/TaskRunner;TaskLifecycleInitializer(98, depends_on=api_wiring) 注入 Task lifecycle directory;InboxInitializer(100, depends_on=api_wiring) 创建或复用 inbox store 并暴露 Inbox directory;随后 SchedulerInitializer(110)ConversationInitializer(120)ChannelInitializer(129) 消费这些 directory。

关键链路: