首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >WorkBuddy 企业实战:构建 AI 数字同事的完整架构与代码实现

WorkBuddy 企业实战:构建 AI 数字同事的完整架构与代码实现

原创
作者头像
资源大佬 jzit-top
发布2026-09-22 15:59:26
发布2026-09-22 15:59:26
60
举报

摘要:WorkBuddy 是企业内部的 AI 数字同事(Digital Coworker),它不只是聊天机器人,而是能理解业务上下文、调用企业系统、执行多步任务、主动提醒、持续学习的智能体。本文从专业视角拆解 WorkBuddy 的核心能力模型与架构,并给出一套可运行的 Python 实现:覆盖 Agent 循环、工具注册、记忆系统、RAG 知识库、权限控制、工作流引擎、主动任务调度、MCP 集成、审计与可观测性。


目录

  1. WorkBuddy 是什么:从聊天机器人到数字同事
  2. WorkBuddy 的六大核心能力
  3. 企业级 WorkBuddy 架构设计
  4. 关键技术:Agent 循环、工具调用、记忆、RAG、MCP
  5. 代码实战:构建企业 WorkBuddy
  6. 模块一:数据模型
  7. 模块二:工具注册与调用
  8. 模块三:记忆系统
  9. 模块四:RAG 知识库
  10. 模块五:Agent 核心循环
  11. 模块六:权限与审计
  12. 模块七:工作流与主动任务
  13. 模块八:API 与接入
  14. 模块九:测试
  15. 部署与生产化清单
  16. 常见陷阱
  17. 总结

一、WorkBuddy 是什么:从聊天机器人到数字同事

企业内部的 AI 助手经历了三个阶段:

代码语言:javascript
复制
第一代:FAQ 机器人     -> 关键词匹配,答不上就转人工
第二代:LLM 聊天助手   -> 能对话,但只能"说",不能"做"
第三代:WorkBuddy      -> 能理解、能规划、能调用系统、能执行、能主动提醒

WorkBuddy 的本质是一个嵌入企业工作流的智能体。它具备:

维度

聊天助手

WorkBuddy

交互

被动问答

主动协作

能力

生成文本

调用工具、执行任务

上下文

当前对话

用户、项目、历史、组织知识

记忆

短期 + 长期 + 组织记忆

权限

RBAC + 数据边界

审计

全量可追溯

集成

CRM / ERP / OA / 工单 / 日历

一句话:WorkBuddy = LLM + 工具 + 记忆 + 规划 + 权限 + 审计 + 企业集成。


二、WorkBuddy 的六大核心能力

代码语言:javascript
复制
┌────────────────────────────────────────────┐
│ 1. 理解    理解自然语言、业务术语、上下文    │
│ 2. 规划    把模糊目标拆成可执行步骤          │
│ 3. 执行    调用企业系统 API / 工具           │
│ 4. 记忆    记住用户偏好、项目历史、组织知识  │
│ 5. 协作    与多人、多系统、多 WorkBuddy 协同 │
│ 6. 主动    定时任务、异常提醒、主动建议      │
└────────────────────────────────────────────┘

典型场景:

  • "帮我把本周所有未回复的客户邮件整理成待办清单"
  • "查一下上周 A 项目的进度,列出风险点"
  • "安排一个明天下午和产品团队的评审会"
  • "把这份合同的关键条款提取出来,对比标准模板"
  • "每天上午 9 点给我一份业务看板摘要"

这些都不是"聊天",而是任务执行


三、企业级 WorkBuddy 架构设计

代码语言:javascript
复制
┌───────────────────────────────────────────────────────┐
│                      接入层                            │
│   Web / 移动端 / 飞书 / 钉钉 / Slack / API / Webhook    │
├───────────────────────────────────────────────────────┤
│                      身份层                            │
│   SSO / OIDC / SAML / 租户 / 角色 / 数据边界            │
├───────────────────────────────────────────────────────┤
│                      Agent 层                          │
│   规划器 / 执行器 / 反思器 / 工具选择 / 多智能体协作      │
├───────────────────────────────────────────────────────┤
│                      能力层                            │
│   工具注册表 / RAG / 记忆 / 工作流 / 定时任务            │
├───────────────────────────────────────────────────────┤
│                      集成层                            │
│   MCP / CRM / ERP / OA / 日历 / 邮件 / 工单 / 数据库     │
├───────────────────────────────────────────────────────┤
│                      数据层                            │
│   PostgreSQL / Redis / 向量库 / 对象存储 / KMS           │
├───────────────────────────────────────────────────────┤
│                      治理层                            │
│   权限 / 审计 / 限流 / 内容安全 / 成本 / 可观测性         │
└───────────────────────────────────────────────────────┘

设计原则:

  1. Agent 无状态:状态放数据库和缓存,方便水平扩展;
  2. 工具即插件:工具注册表统一管理,动态启用;
  3. 记忆分层:短期、长期、组织记忆分开存储;
  4. 权限前置:每次工具调用前做权限检查;
  5. 全量审计:每次思考、调用、输出都可追溯;
  6. 人在回路:高风险动作必须人工确认。

四、关键技术:Agent 循环、工具调用、记忆、RAG、MCP

4.1 Agent 循环

代码语言:javascript
复制
用户目标
  -> 理解意图
  -> 检索上下文(记忆 + RAG)
  -> 制定计划
  -> 选择工具
  -> 执行工具
  -> 观察结果
  -> 反思:是否完成?
  -> 未完成:继续循环
  -> 完成:输出结果 + 更新记忆

4.2 工具调用

工具是 Agent 的"手"。每个工具需要:

  • 名称与描述;
  • 参数 Schema;
  • 权限要求;
  • 执行函数;
  • 错误处理;
  • 审计埋点。

4.3 记忆

类型

说明

存储

短期记忆

当前会话上下文

Redis / 内存

长期记忆

用户偏好、历史任务

PostgreSQL

组织记忆

制度、流程、知识

向量库 + 文档库

4.4 RAG

  • 文档切分;
  • 向量化;
  • 检索 TopK;
  • 重排;
  • 注入 Prompt。

4.5 MCP

Model Context Protocol 让 WorkBuddy 可以统一接入外部工具和数据源,无需为每个系统写适配器。


五、代码实战:构建企业 WorkBuddy

5.1 技术栈

  • Python 3.12
  • FastAPI
  • SQLModel
  • SQLite(示例)/ PostgreSQL(生产)
  • Pydantic(工具 Schema)
  • Pytest

5.2 项目结构

代码语言:javascript
复制
workbuddy/
├── app/
│   ├── __init__.py
│   ├── main.py
│   ├── db.py
│   ├── models.py
│   ├── tools.py
│   ├── memory.py
│   ├── rag.py
│   ├── agent.py
│   ├── llm.py
│   ├── security.py
│   ├── audit.py
│   ├── workflow.py
│   └── scheduler.py
├── tests/
│   └── test_workbuddy.py
├── requirements.txt
└── Dockerfile

5.3 依赖

代码语言:javascript
复制
pip install fastapi uvicorn sqlmodel pydantic pytest httpx numpy

requirements.txt

代码语言:javascript
复制
fastapi
uvicorn[standard]
sqlmodel
pydantic
pytest
httpx
numpy

六、模块一:数据模型

app/models.py

代码语言:javascript
复制
from datetime import datetime
from typing import Optional
from sqlmodel import SQLModel, Field


class User(SQLModel, table=True):
    id: Optional[int] = Field(default=None, primary_key=True)
    username: str = Field(index=True, unique=True)
    display_name: str = ""
    department: str = ""
    role: str = "member"              # admin / manager / member / guest
    enabled: bool = True
    created_at: datetime = Field(default_factory=datetime.utcnow)


class Conversation(SQLModel, table=True):
    id: Optional[int] = Field(default=None, primary_key=True)
    user_id: int = Field(index=True)
    title: str = ""
    created_at: datetime = Field(default_factory=datetime.utcnow)
    updated_at: datetime = Field(default_factory=datetime.utcnow)


class Message(SQLModel, table=True):
    id: Optional[int] = Field(default=None, primary_key=True)
    conversation_id: int = Field(index=True)
    role: str                          # user / assistant / tool / system
    content: str
    tool_name: Optional[str] = None
    created_at: datetime = Field(default_factory=datetime.utcnow)


class Memory(SQLModel, table=True):
    id: Optional[int] = Field(default=None, primary_key=True)
    user_id: int = Field(index=True)
    kind: str                          # preference / fact / project
    key: str
    value: str
    created_at: datetime = Field(default_factory=datetime.utcnow)


class Document(SQLModel, table=True):
    id: Optional[int] = Field(default=None, primary_key=True)
    title: str
    content: str
    department: str = ""
    tags: str = ""
    created_at: datetime = Field(default_factory=datetime.utcnow)


class ToolCallLog(SQLModel, table=True):
    id: Optional[int] = Field(default=None, primary_key=True)
    user_id: int
    conversation_id: int
    tool_name: str
    arguments: str
    result: str
    status: str = "ok"                 # ok / error / denied
    created_at: datetime = Field(default_factory=datetime.utcnow)


class ScheduledTask(SQLModel, table=True):
    id: Optional[int] = Field(default=None, primary_key=True)
    user_id: int = Field(index=True)
    name: str
    cron: str
    prompt: str
    enabled: bool = True
    last_run_at: Optional[datetime] = None
    created_at: datetime = Field(default_factory=datetime.utcnow)

七、模块二:工具注册与调用

7.1 工具基类

app/tools.py

代码语言:javascript
复制
from dataclasses import dataclass, field
from typing import Any, Callable, Dict, List
from pydantic import BaseModel


@dataclass
class ToolSpec:
    name: str
    description: str
    parameters: Dict[str, Any]
    handler: Callable[..., Any]
    requires_roles: List[str] = field(default_factory=lambda: ["member"])
    risk_level: str = "low"           # low / medium / high
    requires_confirmation: bool = False


class ToolRegistry:
    def __init__(self):
        self._tools: Dict[str, ToolSpec] = {}

    def register(self, spec: ToolSpec):
        self._tools[spec.name] = spec

    def get(self, name: str) -> ToolSpec | None:
        return self._tools.get(name)

    def list_for_role(self, role: str) -> List[ToolSpec]:
        return [
            t for t in self._tools.values()
            if role in t.requires_roles
        ]

    def to_llm_schema(self, role: str) -> List[dict]:
        return [
            {
                "name": t.name,
                "description": t.description,
                "parameters": t.parameters,
            }
            for t in self.list_for_role(role)
        ]


registry = ToolRegistry()

7.2 企业工具示例

代码语言:javascript
复制
def tool_calendar_create(user_id: int, title: str, start: str, duration_min: int = 60) -> dict:
    # 生产环境接入真实日历 API
    return {
        "event_id": f"evt_{user_id}_{abs(hash(title)) % 100000}",
        "title": title,
        "start": start,
        "duration_min": duration_min,
        "status": "created",
    }


def tool_ticket_create(user_id: int, title: str, priority: str = "medium") -> dict:
    return {
        "ticket_id": f"TK{abs(hash(title)) % 100000}",
        "title": title,
        "priority": priority,
        "status": "open",
    }


def tool_email_list_unread(user_id: int, limit: int = 10) -> dict:
    # 生产环境接入邮件系统
    return {
        "count": 3,
        "messages": [
            {"from": "client-a@example.com", "subject": "合同确认", "priority": "high"},
            {"from": "partner@example.com", "subject": "合作邀约", "priority": "medium"},
            {"from": "internal@example.com", "subject": "周报提醒", "priority": "low"},
        ][:limit],
    }


def tool_project_status(project_name: str) -> dict:
    return {
        "project": project_name,
        "progress": 0.62,
        "risks": ["关键路径延迟 3 天", "客户 SSO 对接未完成"],
        "owner": "alice",
    }


def tool_kb_search(query: str, top_k: int = 3) -> dict:
    from app.rag import search_documents
    return {"results": search_documents(query, top_k=top_k)}


def register_default_tools():
    registry.register(ToolSpec(
        name="calendar_create",
        description="创建一个日历事件",
        parameters={
            "type": "object",
            "properties": {
                "title": {"type": "string"},
                "start": {"type": "string", "description": "ISO8601 时间"},
                "duration_min": {"type": "integer", "default": 60},
            },
            "required": ["title", "start"],
        },
        handler=tool_calendar_create,
        requires_roles=["member", "manager", "admin"],
        risk_level="medium",
    ))

    registry.register(ToolSpec(
        name="ticket_create",
        description="创建一个工单",
        parameters={
            "type": "object",
            "properties": {
                "title": {"type": "string"},
                "priority": {"type": "string", "enum": ["low", "medium", "high", "urgent"]},
            },
            "required": ["title"],
        },
        handler=tool_ticket_create,
        requires_roles=["member", "manager", "admin"],
        risk_level="low",
    ))

    registry.register(ToolSpec(
        name="email_list_unread",
        description="列出当前用户未读邮件",
        parameters={
            "type": "object",
            "properties": {
                "limit": {"type": "integer", "default": 10},
            },
        },
        handler=tool_email_list_unread,
        requires_roles=["member", "manager", "admin"],
        risk_level="low",
    ))

    registry.register(ToolSpec(
        name="project_status",
        description="查询项目进度与风险",
        parameters={
            "type": "object",
            "properties": {
                "project_name": {"type": "string"},
            },
            "required": ["project_name"],
        },
        handler=tool_project_status,
        requires_roles=["member", "manager", "admin"],
        risk_level="low",
    ))

    registry.register(ToolSpec(
        name="kb_search",
        description="在企业知识库中检索文档",
        parameters={
            "type": "object",
            "properties": {
                "query": {"type": "string"},
                "top_k": {"type": "integer", "default": 3},
            },
            "required": ["query"],
        },
        handler=tool_kb_search,
        requires_roles=["member", "manager", "admin"],
        risk_level="low",
    ))


register_default_tools()

7.3 工具执行器(含权限与审计)

代码语言:javascript
复制
import json
from sqlmodel import Session
from app.models import ToolCallLog
from app.tools import registry, ToolSpec


class ToolExecutor:
    def __init__(self, db: Session, user, conversation_id: int):
        self.db = db
        self.user = user
        self.conversation_id = conversation_id

    def execute(self, tool_name: str, arguments: dict) -> dict:
        spec: ToolSpec | None = registry.get(tool_name)
        if not spec:
            return {"error": f"未知工具:{tool_name}"}

        if self.user.role not in spec.requires_roles:
            self._log(tool_name, arguments, "permission denied", "denied")
            return {"error": "无权限调用该工具"}

        try:
            result = spec.handler(user_id=self.user.id, **arguments)
            self._log(tool_name, arguments, result, "ok")
            return result
        except TypeError as e:
            self._log(tool_name, arguments, str(e), "error")
            return {"error": f"参数错误:{e}"}
        except Exception as e:
            self._log(tool_name, arguments, str(e), "error")
            return {"error": f"执行失败:{e}"}

    def _log(self, tool_name: str, args: dict, result, status: str):
        log = ToolCallLog(
            user_id=self.user.id,
            conversation_id=self.conversation_id,
            tool_name=tool_name,
            arguments=json.dumps(args, ensure_ascii=False),
            result=json.dumps(result, ensure_ascii=False, default=str),
            status=status,
        )
        self.db.add(log)
        self.db.commit()

八、模块三:记忆系统

app/memory.py

代码语言:javascript
复制
from sqlmodel import Session, select
from app.models import Memory


def remember(db: Session, user_id: int, kind: str, key: str, value: str):
    existing = db.exec(
        select(Memory).where(
            Memory.user_id == user_id,
            Memory.kind == kind,
            Memory.key == key,
        )
    ).first()

    if existing:
        existing.value = value
        db.add(existing)
    else:
        db.add(Memory(user_id=user_id, kind=kind, key=key, value=value))
    db.commit()


def recall(db: Session, user_id: int, kind: str | None = None) -> dict:
    stmt = select(Memory).where(Memory.user_id == user_id)
    if kind:
        stmt = stmt.where(Memory.kind == kind)

    rows = db.exec(stmt).all()
    return {r.key: r.value for r in rows}


def build_memory_context(db: Session, user_id: int) -> str:
    memories = recall(db, user_id)
    if not memories:
        return ""

    lines = ["已知用户信息:"]
    for k, v in memories.items():
        lines.append(f"- {k}: {v}")
    return "\n".join(lines)

示例用法:

代码语言:javascript
复制
remember(db, user_id=1, kind="preference", key="时区", value="Asia/Shanghai")
remember(db, user_id=1, kind="preference", key="会议时长", value="默认 45 分钟")
remember(db, user_id=1, kind="fact", key="负责项目", value="AI 客服系统")

九、模块四:RAG 知识库

app/rag.py

代码语言:javascript
复制
import numpy as np
from sqlmodel import Session, select
from app.models import Document
from app.db import engine


def simple_embed(text: str, dim: int = 128) -> np.ndarray:
    """
    极简哈希嵌入,仅用于演示。
    生产环境请替换为 OpenAI / BGE / 本地嵌入模型。
    """
    vec = np.zeros(dim, dtype=np.float32)
    for token in text.lower().split():
        h = hash(token) % dim
        vec[h] += 1.0
    norm = np.linalg.norm(vec) + 1e-9
    return vec / norm


def index_document(db: Session, title: str, content: str, department: str = "", tags: str = ""):
    doc = Document(title=title, content=content, department=department, tags=tags)
    db.add(doc)
    db.commit()
    db.refresh(doc)
    return doc


def search_documents(query: str, top_k: int = 3) -> list:
    with Session(engine) as db:
        docs = db.exec(select(Document)).all()
    if not docs:
        return []

    q_vec = simple_embed(query)
    scored = []
    for d in docs:
        d_vec = simple_embed(d.title + " " + d.content)
        score = float(np.dot(q_vec, d_vec))
        scored.append((score, d))

    scored.sort(key=lambda x: x[0], reverse=True)
    return [
        {
            "id": d.id,
            "title": d.title,
            "snippet": d.content[:200],
            "score": round(score, 4),
        }
        for score, d in scored[:top_k]
    ]

生产环境建议:

  • 使用真实嵌入模型;
  • 使用 FAISS / Milvus / pgvector;
  • 引入重排(rerank);
  • 支持元数据过滤(部门、密级);
  • 加入引用溯源。

十、模块五:Agent 核心循环

app/llm.py

代码语言:javascript
复制
from dataclasses import dataclass


@dataclass
class LLMDecision:
    kind: str                    # "tool" / "final"
    tool_name: str | None = None
    arguments: dict | None = None
    content: str | None = None


def mock_llm_decide(
    user_message: str,
    tool_schemas: list,
    history: list,
    memory_context: str,
) -> LLMDecision:
    """
    Mock LLM:用规则模拟决策,便于本地运行。
    生产环境替换为 OpenAI / Claude 的 function calling。
    """
    text = user_message.lower()

    if "未读邮件" in user_message or "邮件" in user_message:
        return LLMDecision(
            kind="tool",
            tool_name="email_list_unread",
            arguments={"limit": 5},
        )

    if "安排" in user_message and "会" in user_message:
        return LLMDecision(
            kind="tool",
            tool_name="calendar_create",
            arguments={
                "title": "团队评审会",
                "start": "2026-10-01T14:00:00",
                "duration_min": 45,
            },
        )

    if "工单" in user_message:
        return LLMDecision(
            kind="tool",
            tool_name="ticket_create",
            arguments={"title": user_message[:40], "priority": "high"},
        )

    if "项目" in user_message and ("进度" in user_message or "风险" in user_message):
        return LLMDecision(
            kind="tool",
            tool_name="project_status",
            arguments={"project_name": "AI 客服系统"},
        )

    if "知识库" in user_message or "制度" in user_message or "流程" in user_message:
        return LLMDecision(
            kind="tool",
            tool_name="kb_search",
            arguments={"query": user_message, "top_k": 3},
        )

    return LLMDecision(
        kind="final",
        content=f"收到你的消息:{user_message}。我可以帮你查邮件、排会议、查项目、建工单或检索知识库。",
    )

app/agent.py

代码语言:javascript
复制
from sqlmodel import Session
from app.models import Conversation, Message, User
from app.llm import mock_llm_decide, LLMDecision
from app.tools import registry
from app.tools import ToolExecutor
from app.memory import build_memory_context, remember


MAX_STEPS = 5


class WorkBuddyAgent:
    def __init__(self, db: Session, user: User, conversation_id: int):
        self.db = db
        self.user = user
        self.conversation_id = conversation_id
        self.executor = ToolExecutor(db, user, conversation_id)

    def _history(self, limit: int = 20) -> list:
        rows = self.db.query(Message).filter(
            Message.conversation_id == self.conversation_id
        ).order_by(Message.id.desc()).limit(limit).all()
        return list(reversed(rows))

    def _save(self, role: str, content: str, tool_name: str | None = None):
        msg = Message(
            conversation_id=self.conversation_id,
            role=role,
            content=content,
            tool_name=tool_name,
        )
        self.db.add(msg)
        self.db.commit()
        self.db.refresh(msg)
        return msg

    def run(self, user_message: str) -> dict:
        self._save("user", user_message)

        memory_context = build_memory_context(self.db, self.user.id)
        tool_schemas = registry.to_llm_schema(self.user.role)

        steps = []

        for step in range(MAX_STEPS):
            decision: LLMDecision = mock_llm_decide(
                user_message=user_message,
                tool_schemas=tool_schemas,
                history=self._history(),
                memory_context=memory_context,
            )

            if decision.kind == "final":
                self._save("assistant", decision.content or "")
                return {
                    "reply": decision.content,
                    "steps": steps,
                    "conversation_id": self.conversation_id,
                }

            tool_name = decision.tool_name or ""
            args = decision.arguments or {}
            steps.append({"tool": tool_name, "arguments": args})

            result = self.executor.execute(tool_name, args)
            self._save("tool", str(result), tool_name=tool_name)

            # 简单反射:把工具结果拼进下一轮用户消息
            user_message = f"{user_message}\n[工具 {tool_name} 返回] {result}"

            # 记忆沉淀
            if tool_name == "calendar_create":
                remember(
                    self.db, self.user.id,
                    kind="fact", key="最近创建的会议",
                    value=result.get("title", ""),
                )

        final = "任务步骤过多,已停止。请拆分后再试。"
        self._save("assistant", final)
        return {"reply": final, "steps": steps, "conversation_id": self.conversation_id}

说明:真实 LLM 使用 function calling 时,决策由模型输出,无需手写规则。这里的 Mock 是为了让代码可离线运行。生产替换点:把 mock_llm_decide 换成 openai_client.chat.completions.create(..., tools=...) 即可。


十一、模块六:权限与审计

app/security.py

代码语言:javascript
复制
from fastapi import Header, HTTPException, Depends
from sqlmodel import Session, select
from app.db import get_session
from app.models import User

# 示例:生产环境使用 SSO / OIDC
API_KEYS = {
    "admin-key": "admin",
    "manager-key": "manager",
    "member-key": "member",
}


def get_current_user(
    x_api_key: str = Header(..., alias="X-API-Key"),
    session: Session = Depends(get_session),
) -> User:
    role = API_KEYS.get(x_api_key)
    if not role:
        raise HTTPException(status_code=401, detail="Invalid API Key")

    user = session.exec(
        select(User).where(User.username == role)
    ).first()
    if not user:
        raise HTTPException(status_code=401, detail="User not found")
    return user


def require_role(*roles):
    def checker(user: User = Depends(get_current_user)):
        if user.role not in roles:
            raise HTTPException(status_code=403, detail="Forbidden")
        return user
    return checker

app/audit.py

代码语言:javascript
复制
from sqlmodel import Session, select
from app.models import ToolCallLog


def list_tool_calls(db: Session, user_id: int | None = None, limit: int = 200):
    stmt = select(ToolCallLog)
    if user_id:
        stmt = stmt.where(ToolCallLog.user_id == user_id)
    return db.exec(stmt.order_by(ToolCallLog.id.desc()).limit(limit)).all()

十二、模块七:工作流与主动任务

app/workflow.py

代码语言:javascript
复制
from sqlmodel import Session, select
from app.models import ScheduledTask, User
from app.agent import WorkBuddyAgent
from app.models import Conversation


def run_scheduled_task(db: Session, task: ScheduledTask) -> dict:
    user = db.get(User, task.user_id)
    if not user:
        return {"error": "user not found"}

    conv = Conversation(user_id=user.id, title=f"[定时] {task.name}")
    db.add(conv)
    db.commit()
    db.refresh(conv)

    agent = WorkBuddyAgent(db, user, conv.id)
    return agent.run(task.prompt)


def list_tasks(db: Session, user_id: int):
    return db.exec(
        select(ScheduledTask).where(ScheduledTask.user_id == user_id)
    ).all()

app/scheduler.py

代码语言:javascript
复制
from datetime import datetime
from sqlmodel import Session, select
from app.models import ScheduledTask
from app.workflow import run_scheduled_task


def run_due_tasks(db: Session):
    """
    简化版调度:外部 cron 每分钟调用一次。
    生产环境用 APScheduler / Celery Beat / K8s CronJob。
    """
    tasks = db.exec(
        select(ScheduledTask).where(ScheduledTask.enabled == True)  # noqa: E712
    ).all()

    results = []
    for t in tasks:
        try:
            result = run_scheduled_task(db, t)
            t.last_run_at = datetime.utcnow()
            db.add(t)
            db.commit()
            results.append({"task_id": t.id, "result": result})
        except Exception as e:
            results.append({"task_id": t.id, "error": str(e)})
    return results

十三、模块八:API 与接入

app/main.py

代码语言:javascript
复制
from fastapi import FastAPI, Depends, HTTPException
from pydantic import BaseModel
from sqlmodel import Session

from app.db import init_db, get_session
from app.models import User, Conversation, Document, ScheduledTask
from app.security import get_current_user, require_role
from app.agent import WorkBuddyAgent
from app.audit import list_tool_calls
from app.rag import index_document
from app.tools import registry
from app.workflow import list_tasks

app = FastAPI(title="WorkBuddy - Enterprise Digital Coworker")


@app.on_event("startup")
def on_startup():
    init_db()


class ChatRequest(BaseModel):
    message: str
    conversation_id: int | None = None


class DocRequest(BaseModel):
    title: str
    content: str
    department: str = ""
    tags: str = ""


class TaskRequest(BaseModel):
    name: str
    cron: str
    prompt: str


@app.get("/health")
def health():
    return {"status": "ok"}


@app.get("/tools")
def list_tools(user: User = Depends(get_current_user)):
    return [
        {
            "name": t.name,
            "description": t.description,
            "risk_level": t.risk_level,
        }
        for t in registry.list_for_role(user.role)
    ]


@app.post("/chat")
def chat(
    payload: ChatRequest,
    user: User = Depends(get_current_user),
    session: Session = Depends(get_session),
):
    if payload.conversation_id:
        conv = session.get(Conversation, payload.conversation_id)
        if not conv or conv.user_id != user.id:
            raise HTTPException(status_code=404, detail="Conversation not found")
    else:
        conv = Conversation(user_id=user.id, title=payload.message[:30])
        session.add(conv)
        session.commit()
        session.refresh(conv)

    agent = WorkBuddyAgent(session, user, conv.id)
    return agent.run(payload.message)


@app.post("/knowledge")
def add_knowledge(
    payload: DocRequest,
    user: User = Depends(require_role("admin", "manager")),
    session: Session = Depends(get_session),
):
    doc = index_document(
        session,
        title=payload.title,
        content=payload.content,
        department=payload.department,
        tags=payload.tags,
    )
    return {"id": doc.id, "title": doc.title}


@app.post("/schedules")
def create_schedule(
    payload: TaskRequest,
    user: User = Depends(get_current_user),
    session: Session = Depends(get_session),
):
    task = ScheduledTask(
        user_id=user.id,
        name=payload.name,
        cron=payload.cron,
        prompt=payload.prompt,
    )
    session.add(task)
    session.commit()
    session.refresh(task)
    return task


@app.get("/schedules")
def get_schedules(
    user: User = Depends(get_current_user),
    session: Session = Depends(get_session),
):
    return list_tasks(session, user.id)


@app.get("/audit/tool-calls")
def audit(
    user: User = Depends(require_role("admin", "manager")),
    session: Session = Depends(get_session),
):
    return list_tool_calls(session)

十四、模块九:测试

tests/test_workbuddy.py

代码语言:javascript
复制
import os
import pytest
from fastapi.testclient import TestClient

os.environ["DATABASE_URL"] = "sqlite:///./test_workbuddy.db"

from app.main import app  # noqa: E402
from app.db import init_db, engine  # noqa: E402
from app.models import User  # noqa: E402
from sqlmodel import Session  # noqa: E402

client = TestClient(app)


def setup_module():
    init_db()
    with Session(engine) as db:
        for username, role in [
            ("admin", "admin"),
            ("manager", "manager"),
            ("member", "member"),
        ]:
            if not db.query(User).filter(User.username == username).first():
                db.add(User(username=username, display_name=username, role=role))
        db.commit()


def test_health():
    assert client.get("/health").json()["status"] == "ok"


def test_tools_list():
    r = client.get("/tools", headers={"X-API-Key": "member-key"})
    assert r.status_code == 200
    names = [t["name"] for t in r.json()]
    assert "email_list_unread" in names


def test_chat_email():
    r = client.post(
        "/chat",
        json={"message": "帮我看下未读邮件"},
        headers={"X-API-Key": "member-key"},
    )
    assert r.status_code == 200
    data = r.json()
    assert len(data["steps"]) >= 1
    assert data["steps"][0]["tool"] == "email_list_unread"


def test_chat_calendar():
    r = client.post(
        "/chat",
        json={"message": "帮我安排一个团队会"},
        headers={"X-API-Key": "member-key"},
    )
    assert r.status_code == 200
    assert r.json()["steps"][0]["tool"] == "calendar_create"


def test_chat_project():
    r = client.post(
        "/chat",
        json={"message": "查一下 AI 客服系统项目进度和风险"},
        headers={"X-API-Key": "manager-key"},
    )
    assert r.status_code == 200
    assert r.json()["steps"][0]["tool"] == "project_status"


def test_knowledge_and_search():
    r = client.post(
        "/knowledge",
        json={
            "title": "报销制度",
            "content": "员工出差报销需在 7 个工作日内提交,超过需部门经理审批。",
            "department": "财务",
        },
        headers={"X-API-Key": "admin-key"},
    )
    assert r.status_code == 200

    r = client.post(
        "/chat",
        json={"message": "查知识库报销流程"},
        headers={"X-API-Key": "member-key"},
    )
    assert r.status_code == 200
    assert r.json()["steps"][0]["tool"] == "kb_search"


def test_member_cannot_add_knowledge():
    r = client.post(
        "/knowledge",
        json={"title": "x", "content": "y"},
        headers={"X-API-Key": "member-key"},
    )
    assert r.status_code == 403


def test_schedule_create():
    r = client.post(
        "/schedules",
        json={
            "name": "每日业务摘要",
            "cron": "0 9 * * *",
            "prompt": "帮我总结未读邮件和项目风险",
        },
        headers={"X-API-Key": "member-key"},
    )
    assert r.status_code == 200
    assert r.json()["name"] == "每日业务摘要"

运行:

代码语言:javascript
复制
pytest -v

十五、部署与生产化清单

Dockerfile

代码语言:javascript
复制
FROM python:3.12-slim

WORKDIR /app

COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

COPY app ./app

EXPOSE 8000

CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]

启动:

代码语言:javascript
复制
docker build -t workbuddy .
docker run -p 8000:8000 workbuddy

生产化清单:

  • □ 接入真实 LLM(OpenAI / Claude / 本地)
  • □ 接入真实嵌入模型 + 向量库
  • □ SSO / OIDC 替换 API Key
  • □ 工具权限与数据边界(行级)
  • □ 高风险工具二次确认
  • □ 全量审计(思考、调用、输出)
  • □ 内容安全过滤(输入 + 输出)
  • □ 成本监控(Token、调用次数)
  • □ 限流与配额
  • □ 记忆治理(过期、脱敏、可删除)
  • □ RAG 引用溯源
  • □ 调度器换 APScheduler / Celery Beat
  • □ MCP 接入外部系统
  • □ 可观测性(日志、指标、链路追踪)
  • □ 灰度发布与回滚
  • □ 人在回路(Human-in-the-loop)
  • □ 用户反馈闭环
  • □ 数据合规(个人信息、跨境)

十六、常见陷阱

  1. 把 WorkBuddy 做成聊天框:没有工具,就只是玩具。
  2. 工具权限不隔离:低权限用户能调用高权限工具。
  3. 记忆无限增长:成本爆炸,还泄漏隐私。
  4. RAG 没有溯源:用户无法验证答案。
  5. 缺少人在回路:高风险操作直接执行。
  6. 不做审计:出了事无法追责。
  7. Prompt 硬编码:业务变化就要改代码。
  8. 单一 LLM 依赖:模型故障业务中断。
  9. 忽略成本:Token 成本失控。
  10. 不做用户反馈:产品无法迭代。

十七、总结

WorkBuddy 的本质不是"更聪明的聊天机器人",而是嵌入企业工作流的执行型智能体

核心公式:

代码语言:javascript
复制
WorkBuddy = LLM + 工具 + 记忆 + 规划 + 权限 + 审计 + 企业集成

技术主线:

代码语言:javascript
复制
理解 -> 规划 -> 调用工具 -> 观察 -> 反思 -> 输出 -> 记忆沉淀 -> 主动提醒

三条工程原则:

  1. 能力可插拔:工具注册表 + MCP;
  2. 风险可控制:权限 + 确认 + 审计;
  3. 价值可衡量:节省时间、减少错误、提升协作。

代码只是骨架。真正决定 WorkBuddy 成败的,是:

  • 是否真正嵌入业务流程;
  • 是否让用户信任;
  • 是否可持续迭代;
  • 是否守住安全与合规底线。

当 WorkBuddy 能稳定完成"帮我查邮件、排会议、建工单、查项目、检索知识库"这些日常任务时,它就不再是玩具,而是真正的数字同事。


附:完整项目结构

代码语言:javascript
复制
workbuddy/
├── app/
│   ├── __init__.py
│   ├── main.py
│   ├── db.py
│   ├── models.py
│   ├── tools.py
│   ├── memory.py
│   ├── rag.py
│   ├── agent.py
│   ├── llm.py
│   ├── security.py
│   ├── audit.py
│   ├── workflow.py
│   └── scheduler.py
├── tests/
│   └── test_workbuddy.py
├── requirements.txt
└── Dockerfile

运行顺序:

代码语言:javascript
复制
pip install -r requirements.txt
pytest -v
uvicorn app.main:app --reload

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 目录
  • 一、WorkBuddy 是什么:从聊天机器人到数字同事
  • 二、WorkBuddy 的六大核心能力
  • 三、企业级 WorkBuddy 架构设计
  • 四、关键技术:Agent 循环、工具调用、记忆、RAG、MCP
    • 4.1 Agent 循环
    • 4.2 工具调用
    • 4.3 记忆
    • 4.4 RAG
    • 4.5 MCP
  • 五、代码实战:构建企业 WorkBuddy
    • 5.1 技术栈
    • 5.2 项目结构
    • 5.3 依赖
  • 六、模块一:数据模型
  • 七、模块二:工具注册与调用
    • 7.1 工具基类
    • 7.2 企业工具示例
    • 7.3 工具执行器(含权限与审计)
  • 八、模块三:记忆系统
  • 九、模块四:RAG 知识库
  • 十、模块五:Agent 核心循环
  • 十一、模块六:权限与审计
  • 十二、模块七:工作流与主动任务
  • 十三、模块八:API 与接入
  • 十四、模块九:测试
  • 十五、部署与生产化清单
  • 十六、常见陷阱
  • 十七、总结
  • 附:完整项目结构
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档