From d82116d637628f6ae4600442573b935f2d3142c7 Mon Sep 17 00:00:00 2001 From: JiuContinent Date: Wed, 24 Jun 2026 10:30:59 +0800 Subject: [PATCH] =?UTF-8?q?```=20feat(ai=5Fagent):=20=E6=96=B0=E5=A2=9Eope?= =?UTF-8?q?nclaw=5Fhermes=E6=B7=B7=E5=90=88AI=E9=80=82=E9=85=8D=E5=99=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增OpenClawHermesAdapter适配器,结合Hermes记忆功能和OpenClaw执行能力, 实现AI问答流程中的记忆召回、执行操作和记忆存储的完整闭环。 同时更新NoopAdapter提示信息,添加新的模型提供商选项。 feat(business): 新增考勤、工作报告和风险事件业务模型 新增AttendanceRecord、WorkReport、RiskEvent和LegacySyncRun四个业务模型, 扩展业务领域注册表,支持考勤管理、工作报告生成和风险事件跟踪等核心业务功能。 feat(reports): 实现考勤汇总和工作日报周报生成功能 新增attendance_summary方法用于统计每日考勤情况, 新增generate_work_report方法用于生成日/周经营报告, 包含任务完成情况、待处理事项和风险指标等综合信息。 feat(risk): 扩展风险管理API端点和供应商风险检测 新增供应商风险查询端点和风险事件管理端点, 提供风险事件列表查询和自动生成功能, 增强供应商风险评估能力。 feat(feishu): 添加考勤查询命令和风险摘要增强 集成考勤汇总查询功能到飞书命令系统, 在风险摘要中添加供应商风险和开放风险事件统计, 丰富日常经营管理信息展示。 refactor(service): 优化业务服务数据验证和类型转换 重构_model_payload函数实现数据验证和类型转换, 添加列值类型强制转换逻辑,提高API数据处理的准确性和安全性。 build(deps): 添加postgresql数据库驱动依赖 在Dockerfile中添加psycopg[binary]==3.2.3依赖包, 支持PostgreSQL数据库连接和操作。 chore(config): 更新.gitignore文件排除备份和迁移目录 在.gitignore中添加AGENTS.md.bak-*和migration/目录排除规则, 避免备份文件和本地迁移工作区被提交到版本控制系统。 ``` --- .gitignore | 4 + Dockerfile | 1 + app/modules/ai_agent/adapters.py | 83 ++++++++- app/modules/business/models.py | 77 +++++++- app/modules/business/registry.py | 18 +- app/modules/business/service.py | 48 ++++- app/modules/feishu/commands.py | 21 +++ app/modules/legacy_mysql/schemas.py | 1 + app/modules/legacy_mysql/service.py | 18 +- app/modules/reports/routes.py | 33 +++- app/modules/reports/schemas.py | 14 ++ app/modules/reports/service.py | 266 +++++++++++++++++++++++++--- app/modules/risk/routes.py | 19 ++ app/modules/risk/service.py | 230 +++++++++++++++++++++++- app/tools/init_db.py | 8 + company_ai_local.db | Bin 303104 -> 0 bytes docker-compose.yml | 24 --- environment.yml | 1 + logs/feishu_long_connection.err.log | 28 +-- logs/feishu_long_connection.out.log | 2 +- tests/test_ai_adapters.py | 72 ++++++++ tests/test_smoke.py | 66 +++++++ 22 files changed, 953 insertions(+), 81 deletions(-) delete mode 100644 company_ai_local.db create mode 100644 tests/test_ai_adapters.py diff --git a/.gitignore b/.gitignore index 08c54ff..b69ea28 100644 --- a/.gitignore +++ b/.gitignore @@ -9,3 +9,7 @@ __pycache__/ /read.md /README.md /AGENTS.md +/AGENTS.md.bak-* + +# Local migration workspace +/migration/ diff --git a/Dockerfile b/Dockerfile index d2c8619..e3e1101 100644 --- a/Dockerfile +++ b/Dockerfile @@ -11,6 +11,7 @@ RUN pip install --no-cache-dir \ "uvicorn[standard]==0.34.0" \ sqlalchemy==2.0.36 \ pymysql==1.1.1 \ + "psycopg[binary]==3.2.3" \ pydantic-settings==2.7.1 \ python-dotenv==1.0.1 \ httpx==0.28.1 \ diff --git a/app/modules/ai_agent/adapters.py b/app/modules/ai_agent/adapters.py index 997d12c..a969e29 100644 --- a/app/modules/ai_agent/adapters.py +++ b/app/modules/ai_agent/adapters.py @@ -26,7 +26,8 @@ class NoopAdapter(AIAdapter): return { "answer": ( "AI provider is not configured yet. This is a deterministic placeholder. " - "Set MODEL_PROVIDER to openclaw, hermes, or direct_llm after credentials are ready." + "Set MODEL_PROVIDER to openclaw_hermes, openclaw, hermes, or direct_llm " + "after credentials are ready." ), "raw": {"prompt": prompt, "context": context or {}}, } @@ -76,6 +77,78 @@ class HermesAdapter(AIAdapter): return {"answer": data.get("answer") or data.get("content") or str(data), "raw": data} +class OpenClawHermesAdapter(AIAdapter): + """Compose Hermes memory with OpenClaw execution.""" + + provider_name = "openclaw_hermes" + + def __init__(self, settings: Settings): + self.openclaw = OpenClawAdapter(settings) + self.hermes = HermesAdapter(settings) + + def ask(self, prompt: str, context: dict[str, Any] | None = None) -> dict[str, Any]: + base_context = context or {} + recall = self._recall_memory(prompt, base_context) + openclaw_context = { + **base_context, + "agent_pipeline": self.provider_name, + "hermes_memory": recall["answer"], + } + openclaw_result = self.openclaw.ask(prompt, openclaw_context) + remember = self._remember_interaction(prompt, base_context, openclaw_result["answer"]) + return { + "answer": openclaw_result["answer"], + "raw": { + "pipeline": "hermes_recall -> openclaw_answer -> hermes_remember", + "hermes_recall": recall, + "openclaw": openclaw_result.get("raw", {}), + "hermes_remember": remember, + }, + } + + def _recall_memory(self, prompt: str, context: dict[str, Any]) -> dict[str, Any]: + recall_prompt = ( + "Retrieve concise long-term memory, preferences, prior decisions, and relevant " + "business context for this request. Return only information useful to answer it." + ) + try: + result = self.hermes.ask( + recall_prompt, + { + "mode": "memory_recall", + "user_prompt": prompt, + "request_context": context, + }, + ) + except Exception as exc: # Hermes memory should not block OpenClaw execution. + return {"answer": "", "raw": {}, "error": _error_detail(exc)} + return {"answer": result["answer"], "raw": result.get("raw", {})} + + def _remember_interaction( + self, + prompt: str, + context: dict[str, Any], + answer: str, + ) -> dict[str, Any]: + remember_prompt = ( + "Store durable lessons from this interaction for future company management " + "assistance. Ignore transient details and do not store secrets." + ) + try: + result = self.hermes.ask( + remember_prompt, + { + "mode": "memory_write", + "user_prompt": prompt, + "request_context": context, + "assistant_answer": answer, + }, + ) + except Exception as exc: + return {"ok": False, "raw": {}, "error": _error_detail(exc)} + return {"ok": True, "raw": result.get("raw", {}), "answer": result["answer"]} + + class DirectLLMAdapter(AIAdapter): """Adapter for OpenAI-compatible chat completions APIs.""" @@ -118,6 +191,14 @@ def get_adapter() -> AIAdapter: return OpenClawAdapter(settings) if provider == "hermes": return HermesAdapter(settings) + if provider in {"openclaw_hermes", "openclaw-hermes", "hybrid"}: + return OpenClawHermesAdapter(settings) if provider == "direct_llm": return DirectLLMAdapter(settings) return NoopAdapter() + + +def _error_detail(exc: Exception) -> Any: + if isinstance(exc, HTTPException): + return exc.detail + return {"type": type(exc).__name__, "message": str(exc)} diff --git a/app/modules/business/models.py b/app/modules/business/models.py index f53ac02..83db0b7 100644 --- a/app/modules/business/models.py +++ b/app/modules/business/models.py @@ -1,7 +1,7 @@ from datetime import date, datetime from decimal import Decimal -from sqlalchemy import Date, DateTime, Integer, Numeric, String, Text +from sqlalchemy import JSON, Date, DateTime, Integer, Numeric, String, Text from sqlalchemy.orm import Mapped, mapped_column from app.core.database import Base @@ -157,3 +157,78 @@ class Supplier(Base, TimestampMixin): price_score: Mapped[Decimal] = mapped_column(Numeric(5, 2), default=0) risk_level: Mapped[str] = mapped_column(String(32), default="low", index=True) blacklist_status: Mapped[str] = mapped_column(String(32), default="normal", index=True) + + +class AttendanceRecord(Base, TimestampMixin): + __tablename__ = "attendance_records" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + code: Mapped[str] = mapped_column(String(64), unique=True, index=True) + employee_name: Mapped[str] = mapped_column(String(128), index=True) + employee_id: Mapped[str | None] = mapped_column(String(128), nullable=True, index=True) + department: Mapped[str | None] = mapped_column(String(128), nullable=True, index=True) + project_code: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True) + work_date: Mapped[date] = mapped_column(Date, index=True) + check_in_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + check_out_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + status: Mapped[str] = mapped_column(String(64), default="正常", index=True) + location: Mapped[str | None] = mapped_column(String(255), nullable=True) + source_system: Mapped[str] = mapped_column(String(64), default="internal") + external_id: Mapped[str | None] = mapped_column(String(128), nullable=True, index=True) + note: Mapped[str | None] = mapped_column(Text, nullable=True) + + +class WorkReport(Base, TimestampMixin): + __tablename__ = "work_reports" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + code: Mapped[str] = mapped_column(String(64), unique=True, index=True) + report_type: Mapped[str] = mapped_column(String(32), index=True) + title: Mapped[str] = mapped_column(String(255), index=True) + reporter: Mapped[str] = mapped_column(String(128), default="system", index=True) + department: Mapped[str | None] = mapped_column(String(128), nullable=True, index=True) + project_code: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True) + period_start: Mapped[date] = mapped_column(Date, index=True) + period_end: Mapped[date] = mapped_column(Date, index=True) + content: Mapped[str] = mapped_column(Text) + metrics: Mapped[dict | None] = mapped_column(JSON, nullable=True) + risk_summary: Mapped[dict | None] = mapped_column(JSON, nullable=True) + status: Mapped[str] = mapped_column(String(64), default="已生成", index=True) + source_system: Mapped[str] = mapped_column(String(64), default="internal") + + +class RiskEvent(Base, TimestampMixin): + __tablename__ = "risk_events" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + code: Mapped[str] = mapped_column(String(64), unique=True, index=True) + title: Mapped[str] = mapped_column(String(255), index=True) + risk_type: Mapped[str] = mapped_column(String(64), index=True) + risk_level: Mapped[str] = mapped_column(String(32), default="medium", index=True) + status: Mapped[str] = mapped_column(String(32), default="open", index=True) + source_domain: Mapped[str] = mapped_column(String(128), index=True) + source_record_id: Mapped[str | None] = mapped_column(String(128), nullable=True, index=True) + project_code: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True) + owner: Mapped[str | None] = mapped_column(String(128), nullable=True, index=True) + detected_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) + due_date: Mapped[date | None] = mapped_column(Date, nullable=True, index=True) + description: Mapped[str | None] = mapped_column(Text, nullable=True) + mitigation: Mapped[str | None] = mapped_column(Text, nullable=True) + evidence: Mapped[dict | None] = mapped_column(JSON, nullable=True) + + +class LegacySyncRun(Base, TimestampMixin): + __tablename__ = "legacy_sync_runs" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + code: Mapped[str] = mapped_column(String(64), unique=True, index=True) + domain: Mapped[str] = mapped_column(String(128), index=True) + source_table: Mapped[str | None] = mapped_column(String(128), nullable=True, index=True) + status: Mapped[str] = mapped_column(String(32), default="running", index=True) + started_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) + finished_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + created_count: Mapped[int] = mapped_column(Integer, default=0) + updated_count: Mapped[int] = mapped_column(Integer, default=0) + skipped_count: Mapped[int] = mapped_column(Integer, default=0) + error_message: Mapped[str | None] = mapped_column(Text, nullable=True) + note: Mapped[str | None] = mapped_column(Text, nullable=True) diff --git a/app/modules/business/registry.py b/app/modules/business/registry.py index 5337fae..0fc48df 100644 --- a/app/modules/business/registry.py +++ b/app/modules/business/registry.py @@ -13,9 +13,25 @@ DOMAIN_MODELS: dict[str, type[DeclarativeMeta]] = { "standards": models.Standard, "performance-metrics": models.PerformanceMetric, "suppliers": models.Supplier, + "attendance-records": models.AttendanceRecord, + "work-reports": models.WorkReport, + "risk-events": models.RiskEvent, + "legacy-sync-runs": models.LegacySyncRun, } -LOW_RISK_DOMAINS = {"projects", "tasks", "procurements", "expenses", "policies", "standards"} +LOW_RISK_DOMAINS = { + "projects", + "tasks", + "procurements", + "expenses", + "policies", + "standards", + "suppliers", + "attendance-records", + "work-reports", + "risk-events", + "legacy-sync-runs", +} HIGH_RISK_DOMAINS = {"fund-accounts", "performance-metrics"} diff --git a/app/modules/business/service.py b/app/modules/business/service.py index eaec554..3fcb840 100644 --- a/app/modules/business/service.py +++ b/app/modules/business/service.py @@ -1,9 +1,13 @@ from datetime import date, datetime -from decimal import Decimal +from decimal import Decimal, InvalidOperation from typing import Any from fastapi import HTTPException, status +from sqlalchemy import Date as SQLDate +from sqlalchemy import DateTime as SQLDateTime +from sqlalchemy import Numeric as SQLNumeric from sqlalchemy import Select, func, select +from sqlalchemy.sql.schema import Column from sqlalchemy.orm import Session from app.modules.audit.schemas import AuditLogCreate @@ -27,6 +31,39 @@ def serialize_model(record: Any) -> dict[str, Any]: return data +def _coerce_column_value(column: Column, value: Any) -> Any: + """Coerce API JSON values into the Python type expected by a SQLAlchemy column.""" + + if value is None: + return None + if isinstance(column.type, SQLDateTime) and isinstance(value, str): + return datetime.fromisoformat(value.replace("Z", "+00:00")) + if isinstance(column.type, SQLDate) and isinstance(value, str): + return date.fromisoformat(value) + if isinstance(column.type, SQLNumeric) and not isinstance(value, Decimal): + return Decimal(str(value)) + return value + + +def _model_payload(model: Any, data: dict[str, Any]) -> dict[str, Any]: + """Filter unknown keys and coerce values according to model column types.""" + + columns = {column.name: column for column in model.__table__.columns if column.name != "id"} + payload: dict[str, Any] = {} + for key, value in data.items(): + column = columns.get(key) + if column is None: + continue + try: + payload[key] = _coerce_column_value(column, value) + except (ValueError, TypeError, InvalidOperation) as exc: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + detail=f"Invalid value for field '{key}'", + ) from exc + return payload + + class BusinessService: """Manage generic CRUD operations across registered business domains.""" @@ -65,8 +102,7 @@ class BusinessService: actor: str = "api", ) -> dict[str, Any]: model = get_domain_model(domain) - allowed = {column.name for column in model.__table__.columns if column.name != "id"} - payload = {key: value for key, value in data.items() if key in allowed} + payload = _model_payload(model, data) record = model(**payload) self.db.add(record) self.db.commit() @@ -116,10 +152,8 @@ class BusinessService: status_code=status.HTTP_404_NOT_FOUND, detail="Record not found", ) - allowed = {column.name for column in model.__table__.columns if column.name != "id"} - for key, value in data.items(): - if key in allowed: - setattr(record, key, value) + for key, value in _model_payload(model, data).items(): + setattr(record, key, value) self.db.commit() self.db.refresh(record) result = serialize_model(record) diff --git a/app/modules/feishu/commands.py b/app/modules/feishu/commands.py index 5b9f2f7..c1bc9e9 100644 --- a/app/modules/feishu/commands.py +++ b/app/modules/feishu/commands.py @@ -108,6 +108,25 @@ class FeishuCommandService: result["provider_response"] = provider_response return result + if any(keyword in command_text for keyword in ["打卡", "考勤", "attendance"]): + report = ReportService(self.db).attendance_summary() + result = { + "command": "attendance_summary", + "reply_type": "card", + "title": report["title"], + "content": report["content"], + "lines": report["lines"], + } + if auto_reply: + provider_response = self._send_card_if_configured( + chat_id, + report["title"], + report["lines"], + actor, + ) + result["provider_response"] = provider_response + return result + if any(keyword in command_text for keyword in ["风险", "预警", "risk"]): summary = RiskService(self.db).summary() lines = [ @@ -117,6 +136,8 @@ class FeishuCommandService: f"- 延期项目:{len(summary['delayed_projects'])}", f"- 超预算项目:{len(summary['over_budget_projects'])}", f"- 资金风险账户:{len(summary['fund_risks'])}", + f"- 供应商风险:{len(summary['supplier_risks'])}", + f"- 打开风险事件:{len(summary['open_events'])}", ] result = { "command": "risk_summary", diff --git a/app/modules/legacy_mysql/schemas.py b/app/modules/legacy_mysql/schemas.py index b8cc084..8008571 100644 --- a/app/modules/legacy_mysql/schemas.py +++ b/app/modules/legacy_mysql/schemas.py @@ -51,4 +51,5 @@ class LegacyProjectSyncResult(BaseModel): created: int updated: int skipped: int + sync_run_code: str | None = None items: list[dict[str, Any]] diff --git a/app/modules/legacy_mysql/service.py b/app/modules/legacy_mysql/service.py index 8f905bd..dce8d57 100644 --- a/app/modules/legacy_mysql/service.py +++ b/app/modules/legacy_mysql/service.py @@ -13,7 +13,7 @@ from app.core.config import get_settings from app.core.database import legacy_engine from app.modules.audit.schemas import AuditLogCreate from app.modules.audit.service import AuditService -from app.modules.business.models import Project +from app.modules.business.models import LegacySyncRun, Project from app.modules.business.service import serialize_model FORBIDDEN_SQL_TOKENS = { @@ -256,6 +256,22 @@ class LegacyMySQLService: "skipped": skipped, "items": items, } + sync_run = LegacySyncRun( + code=f"SYNC-PROJECTS-{datetime.utcnow():%Y%m%d%H%M%S%f}", + domain="projects", + source_table="LEGACY_PROJECT_QUERY", + status="dry_run" if dry_run else "success", + finished_at=datetime.utcnow(), + created_count=created, + updated_count=updated, + skipped_count=skipped, + note="Project sync from readonly legacy MySQL", + ) + self.db.add(sync_run) + self.db.commit() + self.db.refresh(sync_run) + result["sync_run_code"] = sync_run.code + AuditService(self.db).log( AuditLogCreate( actor=actor, diff --git a/app/modules/reports/routes.py b/app/modules/reports/routes.py index b9a1d48..7b48cda 100644 --- a/app/modules/reports/routes.py +++ b/app/modules/reports/routes.py @@ -1,9 +1,15 @@ +from datetime import date + from fastapi import APIRouter, Depends from sqlalchemy.orm import Session from app.core.database import get_db from app.core.security import require_api_key -from app.modules.reports.schemas import PushReportRequest, ReportResponse +from app.modules.reports.schemas import ( + PushReportRequest, + ReportResponse, + WorkReportGenerateRequest, +) from app.modules.reports.service import ReportService router = APIRouter(dependencies=[Depends(require_api_key)]) @@ -19,6 +25,31 @@ def project_weekly(db: Session = Depends(get_db)) -> dict: return ReportService(db).project_weekly() +@router.get("/attendance-summary") +def attendance_summary( + work_date: date | None = None, + db: Session = Depends(get_db), +) -> dict: + return ReportService(db).attendance_summary(work_date) + + +@router.post("/work-reports/generate") +def generate_work_report( + payload: WorkReportGenerateRequest, + db: Session = Depends(get_db), +) -> dict: + return ReportService(db).generate_work_report( + report_type=payload.report_type, + reporter=payload.reporter, + department=payload.department, + project_code=payload.project_code, + period_start=payload.period_start, + period_end=payload.period_end, + persist=payload.persist, + actor=payload.actor, + ) + + @router.post("/daily-brief/push") def push_daily_brief(payload: PushReportRequest, db: Session = Depends(get_db)) -> dict: report = ReportService(db).daily_brief() diff --git a/app/modules/reports/schemas.py b/app/modules/reports/schemas.py index 168481e..a1a1832 100644 --- a/app/modules/reports/schemas.py +++ b/app/modules/reports/schemas.py @@ -1,3 +1,6 @@ +from datetime import date +from typing import Literal + from pydantic import BaseModel @@ -11,3 +14,14 @@ class PushReportRequest(BaseModel): receive_id: str | None = None receive_id_type: str = "chat_id" actor: str = "system" + + +class WorkReportGenerateRequest(BaseModel): + report_type: Literal["daily", "weekly"] = "daily" + reporter: str = "system" + department: str | None = None + project_code: str | None = None + period_start: date | None = None + period_end: date | None = None + persist: bool = True + actor: str = "api" diff --git a/app/modules/reports/service.py b/app/modules/reports/service.py index 2593a4d..8041e77 100644 --- a/app/modules/reports/service.py +++ b/app/modules/reports/service.py @@ -1,12 +1,30 @@ +from datetime import date, datetime, timedelta from decimal import Decimal +from typing import Any from sqlalchemy import func, select from sqlalchemy.orm import Session -from app.modules.business.models import Expense, FundAccount, Procurement, Project, WorkTask +from app.modules.audit.schemas import AuditLogCreate +from app.modules.audit.service import AuditService +from app.modules.business.models import ( + AttendanceRecord, + Expense, + FundAccount, + Procurement, + Project, + RiskEvent, + WorkReport, + WorkTask, +) +from app.modules.business.service import serialize_model from app.modules.feishu.service import FeishuService from app.modules.risk.service import RiskService +DONE_STATUSES = {"完成", "已完成", "关闭", "done", "completed", "closed"} +PENDING_APPROVAL_STATUSES = {"草稿", "审批中", "待审批", "pending"} +PROJECT_CLOSED_STATUSES = {"验收", "已完成", "复盘", "归档", "关闭", "closed"} + def _money(value: Decimal | int | float | None) -> str: """Format a numeric value as a two-decimal money string.""" @@ -15,6 +33,26 @@ def _money(value: Decimal | int | float | None) -> str: return f"{amount:,.2f}" +def _json_safe(value: Any) -> Any: + """Convert nested report payloads into JSON-storable values.""" + + if isinstance(value, Decimal): + return float(value) + if isinstance(value, (datetime, date)): + return value.isoformat() + if isinstance(value, list): + return [_json_safe(item) for item in value] + if isinstance(value, dict): + return {key: _json_safe(item) for key, item in value.items()} + return value + + +def _next_code(prefix: str) -> str: + """Build a compact unique code for generated report records.""" + + return f"{prefix}-{datetime.utcnow():%Y%m%d%H%M%S%f}" + + class ReportService: """Build operational reports and push them through Feishu.""" @@ -22,61 +60,59 @@ class ReportService: self.db = db self.risks = RiskService(db) + def _count(self, model: type, *conditions: Any) -> int: + stmt = select(func.count()).select_from(model) + for condition in conditions: + stmt = stmt.where(condition) + return int(self.db.execute(stmt).scalar() or 0) + def daily_brief(self) -> dict: - project_count = int( - self.db.execute(select(func.count()).select_from(Project)).scalar() or 0 + project_count = self._count(Project) + task_count = self._count(WorkTask) + procurement_pending = self._count( + Procurement, + Procurement.approval_status.in_(PENDING_APPROVAL_STATUSES), ) - task_count = int(self.db.execute(select(func.count()).select_from(WorkTask)).scalar() or 0) - procurement_pending = int( - self.db.execute( - select(func.count()) - .select_from(Procurement) - .where(Procurement.approval_status.in_(["草稿", "审批中", "待审批"])) - ).scalar() - or 0 - ) - expense_pending = int( - self.db.execute( - select(func.count()) - .select_from(Expense) - .where(Expense.approval_status.in_(["草稿", "审批中", "待审批"])) - ).scalar() - or 0 + expense_pending = self._count( + Expense, + Expense.approval_status.in_(PENDING_APPROVAL_STATUSES), ) fund_total = ( self.db.execute(select(func.sum(FundAccount.current_balance))).scalar() or Decimal("0") ) risk_summary = self.risks.summary() + attendance = self.attendance_summary() lines = [ f"- 项目总数:{project_count}", f"- 任务总数:{task_count}", f"- 待处理采购:{procurement_pending}", f"- 待处理费用:{expense_pending}", f"- 当前账户总余额:{_money(fund_total)}", + f"- 今日打卡记录:{attendance['total']},异常:{attendance['abnormal_total']}", f"- 逾期任务:{len(risk_summary['overdue_tasks'])}", f"- 延期项目:{len(risk_summary['delayed_projects'])}", f"- 超预算项目:{len(risk_summary['over_budget_projects'])}", f"- 资金风险账户:{len(risk_summary['fund_risks'])}", + f"- 供应商风险:{len(risk_summary['supplier_risks'])}", + f"- 打开风险事件:{len(risk_summary['open_events'])}", f"- 综合风险等级:{risk_summary['risk_level']}", ] return {"title": "每日经营晨报", "lines": lines, "content": "\n".join(lines)} def project_weekly(self) -> dict: - active = int( - self.db.execute( - select(func.count()) - .select_from(Project) - .where(Project.status.notin_(["验收", "已完成", "复盘", "归档", "关闭"])) - ).scalar() - or 0 + active = self._count( + Project, + Project.status.notin_(PROJECT_CLOSED_STATUSES), ) delayed = self.risks.delayed_projects() over_budget = self.risks.over_budget_projects() + open_risks = self.risks.list_events(status_filter="open") lines = [ f"- 活跃项目:{active}", f"- 延期项目:{len(delayed)}", f"- 超预算项目:{len(over_budget)}", + f"- 打开风险事件:{len(open_risks)}", "- 需要管理层关注:", ] for item in delayed[:10]: @@ -85,6 +121,184 @@ class ReportService: lines.append(f" - 超预算:{item.get('code')} {item.get('name')}") return {"title": "项目周报", "lines": lines, "content": "\n".join(lines)} + def attendance_summary(self, work_date: date | None = None) -> dict[str, Any]: + """Summarize attendance records for one business day.""" + + target_date = work_date or date.today() + rows = self.db.execute( + select(AttendanceRecord.status, func.count()) + .where(AttendanceRecord.work_date == target_date) + .group_by(AttendanceRecord.status) + ).all() + status_counts = {str(status): int(count) for status, count in rows} + abnormal_total = sum( + count + for status, count in status_counts.items() + if status in {"迟到", "早退", "缺卡", "旷工", "异常"} + ) + total = sum(status_counts.values()) + lines = [ + f"- 日期:{target_date.isoformat()}", + f"- 打卡记录:{total}", + f"- 异常记录:{abnormal_total}", + ] + for status, count in sorted(status_counts.items()): + lines.append(f"- {status}:{count}") + return { + "title": "打卡汇总", + "work_date": target_date.isoformat(), + "total": total, + "abnormal_total": abnormal_total, + "status_counts": status_counts, + "lines": lines, + "content": "\n".join(lines), + } + + def generate_work_report( + self, + report_type: str = "daily", + reporter: str = "system", + department: str | None = None, + project_code: str | None = None, + period_start: date | None = None, + period_end: date | None = None, + persist: bool = True, + actor: str = "api", + ) -> dict[str, Any]: + """Generate a daily or weekly operating report, optionally persisting it.""" + + start, end = self._resolve_period(report_type, period_start, period_end) + metrics = self._report_metrics(start, end, project_code, department) + risk_summary = _json_safe(self.risks.summary()) + title = "经营日报" if report_type == "daily" else "经营周报" + lines = self._work_report_lines(title, start, end, metrics, risk_summary) + report = { + "title": title, + "report_type": report_type, + "period_start": start.isoformat(), + "period_end": end.isoformat(), + "lines": lines, + "content": "\n".join(lines), + "metrics": metrics, + "risk_summary": risk_summary, + } + + record_data = None + if persist: + record = WorkReport( + code=_next_code(f"REPORT-{report_type.upper()}"), + report_type=report_type, + title=title, + reporter=reporter, + department=department, + project_code=project_code, + period_start=start, + period_end=end, + content=report["content"], + metrics=metrics, + risk_summary=risk_summary, + ) + self.db.add(record) + self.db.commit() + self.db.refresh(record) + record_data = serialize_model(record) + AuditService(self.db).log( + AuditLogCreate( + actor=actor, + source="reports", + action=f"generate_{report_type}_report", + target_type="work-reports", + target_id=str(record.id), + response_payload=record_data, + ) + ) + + return {"report": report, "data": record_data} + + def _resolve_period( + self, + report_type: str, + period_start: date | None, + period_end: date | None, + ) -> tuple[date, date]: + today = date.today() + if report_type == "daily": + start = period_start or period_end or today + return start, period_end or start + end = period_end or today + start = period_start or end - timedelta(days=6) + return start, end + + def _report_metrics( + self, + start: date, + end: date, + project_code: str | None, + department: str | None, + ) -> dict[str, Any]: + task_filters = [ + WorkTask.due_date.is_not(None), + WorkTask.due_date >= start, + WorkTask.due_date <= end, + ] + procurement_filters = [Procurement.approval_status.in_(PENDING_APPROVAL_STATUSES)] + expense_filters = [Expense.approval_status.in_(PENDING_APPROVAL_STATUSES)] + attendance_filters = [ + AttendanceRecord.work_date >= start, + AttendanceRecord.work_date <= end, + ] + if project_code: + task_filters.append(WorkTask.project_code == project_code) + procurement_filters.append(Procurement.project_code == project_code) + expense_filters.append(Expense.project_code == project_code) + attendance_filters.append(AttendanceRecord.project_code == project_code) + if department: + expense_filters.append(Expense.department == department) + attendance_filters.append(AttendanceRecord.department == department) + + completed_tasks = self._count(WorkTask, WorkTask.status.in_(DONE_STATUSES), *task_filters) + overdue_tasks = self._count( + WorkTask, + WorkTask.due_date < date.today(), + WorkTask.status.notin_(DONE_STATUSES), + *task_filters, + ) + return { + "projects_total": self._count(Project), + "active_projects": self._count( + Project, + Project.status.notin_(PROJECT_CLOSED_STATUSES), + ), + "tasks_total": self._count(WorkTask, *task_filters), + "tasks_completed": completed_tasks, + "tasks_overdue": overdue_tasks, + "procurements_pending": self._count(Procurement, *procurement_filters), + "expenses_pending": self._count(Expense, *expense_filters), + "attendance_total": self._count(AttendanceRecord, *attendance_filters), + "open_risk_events": self._count(RiskEvent, RiskEvent.status == "open"), + } + + def _work_report_lines( + self, + title: str, + start: date, + end: date, + metrics: dict[str, Any], + risk_summary: dict[str, Any], + ) -> list[str]: + return [ + f"- 报告:{title}", + f"- 周期:{start.isoformat()} 至 {end.isoformat()}", + f"- 项目:总数 {metrics['projects_total']},活跃 {metrics['active_projects']}", + f"- 任务:总数 {metrics['tasks_total']},完成 {metrics['tasks_completed']}", + f"- 逾期任务:{metrics['tasks_overdue']}", + f"- 待处理采购:{metrics['procurements_pending']}", + f"- 待处理费用:{metrics['expenses_pending']}", + f"- 打卡记录:{metrics['attendance_total']}", + f"- 打开风险事件:{metrics['open_risk_events']}", + f"- 综合风险等级:{risk_summary['risk_level']}", + ] + def push_report( self, report: dict, diff --git a/app/modules/risk/routes.py b/app/modules/risk/routes.py index 4effb31..02ef471 100644 --- a/app/modules/risk/routes.py +++ b/app/modules/risk/routes.py @@ -31,3 +31,22 @@ def over_budget_projects(db: Session = Depends(get_db)) -> dict: @router.get("/funds") def fund_risks(db: Session = Depends(get_db)) -> dict: return {"items": RiskService(db).fund_risks()} + + +@router.get("/suppliers") +def supplier_risks(db: Session = Depends(get_db)) -> dict: + return {"items": RiskService(db).supplier_risks()} + + +@router.get("/events") +def risk_events( + limit: int = 100, + status: str | None = None, + db: Session = Depends(get_db), +) -> dict: + return {"items": RiskService(db).list_events(limit=limit, status_filter=status)} + + +@router.post("/events/generate") +def generate_risk_events(actor: str = "api", db: Session = Depends(get_db)) -> dict: + return RiskService(db).generate_events(actor=actor) diff --git a/app/modules/risk/service.py b/app/modules/risk/service.py index 51b9979..44f1708 100644 --- a/app/modules/risk/service.py +++ b/app/modules/risk/service.py @@ -1,13 +1,18 @@ -from datetime import date +from datetime import date, datetime from decimal import Decimal from typing import Any from sqlalchemy import select from sqlalchemy.orm import Session -from app.modules.business.models import FundAccount, Project, WorkTask +from app.modules.audit.schemas import AuditLogCreate +from app.modules.audit.service import AuditService +from app.modules.business.models import FundAccount, Project, RiskEvent, Supplier, WorkTask from app.modules.business.service import serialize_model +DONE_STATUSES = {"完成", "已完成", "关闭", "done", "completed", "closed"} +CLOSED_RISK_STATUSES = {"closed", "resolved"} + class RiskService: """Evaluate rule-based business risk signals from internal ledgers.""" @@ -19,7 +24,7 @@ class RiskService: stmt = select(WorkTask).where( WorkTask.due_date.is_not(None), WorkTask.due_date < date.today(), - WorkTask.status.notin_(["完成", "已完成", "关闭"]), + WorkTask.status.notin_(DONE_STATUSES), ) return [serialize_model(item) for item in self.db.execute(stmt).scalars()] @@ -27,7 +32,7 @@ class RiskService: stmt = select(Project).where( Project.due_date.is_not(None), Project.due_date < date.today(), - Project.status.notin_(["验收", "已完成", "复盘", "归档", "关闭"]), + Project.status.notin_(["验收", "已完成", "复盘", "归档", "关闭", "closed"]), ) return [serialize_model(item) for item in self.db.execute(stmt).scalars()] @@ -42,16 +47,42 @@ class RiskService: stmt = select(FundAccount).where(FundAccount.current_balance < FundAccount.safety_line) return [serialize_model(item) for item in self.db.execute(stmt).scalars()] + def supplier_risks(self) -> list[dict[str, Any]]: + stmt = select(Supplier).where( + (Supplier.blacklist_status != "normal") + | Supplier.risk_level.in_(["medium", "high"]) + ) + return [serialize_model(item) for item in self.db.execute(stmt).scalars()] + + def list_events( + self, + limit: int = 100, + status_filter: str | None = None, + ) -> list[dict[str, Any]]: + stmt = select(RiskEvent).order_by(RiskEvent.id.desc()).limit(min(limit, 500)) + if status_filter: + stmt = ( + select(RiskEvent) + .where(RiskEvent.status == status_filter) + .order_by(RiskEvent.id.desc()) + .limit(min(limit, 500)) + ) + return [serialize_model(item) for item in self.db.execute(stmt).scalars()] + def summary(self) -> dict[str, Any]: overdue_tasks = self.overdue_tasks() delayed_projects = self.delayed_projects() over_budget_projects = self.over_budget_projects() fund_risks = self.fund_risks() + supplier_risks = self.supplier_risks() + open_events = self.list_events(status_filter="open") risk_score = ( len(overdue_tasks) * 1 + len(delayed_projects) * 3 + len(over_budget_projects) * 4 + len(fund_risks) * 5 + + len(supplier_risks) * 3 + + len(open_events) * 2 ) if risk_score >= 15: level = "high" @@ -66,4 +97,195 @@ class RiskService: "delayed_projects": delayed_projects, "over_budget_projects": over_budget_projects, "fund_risks": fund_risks, + "supplier_risks": supplier_risks, + "open_events": open_events, } + + def generate_events(self, actor: str = "api") -> dict[str, Any]: + """Generate or refresh risk-event ledger entries from current signals.""" + + payloads = self._build_event_payloads() + created = 0 + updated = 0 + skipped = 0 + items: list[dict[str, Any]] = [] + + for payload in payloads: + record = self.db.execute( + select(RiskEvent).where(RiskEvent.code == payload["code"]) + ).scalar_one_or_none() + if record is None: + record = RiskEvent(**payload) + self.db.add(record) + self.db.flush() + created += 1 + action = "created" + elif record.status in CLOSED_RISK_STATUSES: + skipped += 1 + items.append({"action": "skipped", "risk_event": serialize_model(record)}) + continue + else: + for key, value in payload.items(): + if key != "code": + setattr(record, key, value) + updated += 1 + action = "updated" + items.append({"action": action, "risk_event": serialize_model(record)}) + + self.db.commit() + AuditService(self.db).log( + AuditLogCreate( + actor=actor, + source="risk", + action="generate_events", + target_type="risk-events", + risk_level="medium", + response_payload={ + "created": created, + "updated": updated, + "skipped": skipped, + }, + ) + ) + return {"created": created, "updated": updated, "skipped": skipped, "items": items} + + def _build_event_payloads(self) -> list[dict[str, Any]]: + payloads: list[dict[str, Any]] = [] + payloads.extend(self._overdue_task_payloads()) + payloads.extend(self._delayed_project_payloads()) + payloads.extend(self._over_budget_project_payloads()) + payloads.extend(self._fund_risk_payloads()) + payloads.extend(self._supplier_risk_payloads()) + return payloads + + def _overdue_task_payloads(self) -> list[dict[str, Any]]: + stmt = select(WorkTask).where( + WorkTask.due_date.is_not(None), + WorkTask.due_date < date.today(), + WorkTask.status.notin_(DONE_STATUSES), + ) + payloads = [] + for task in self.db.execute(stmt).scalars(): + payloads.append( + { + "code": f"RISK-TASK-OVERDUE-{task.id}", + "title": f"任务逾期:{task.title}", + "risk_type": "overdue_task", + "risk_level": "medium", + "status": "open", + "source_domain": "tasks", + "source_record_id": str(task.id), + "project_code": task.project_code, + "owner": task.owner, + "due_date": task.due_date, + "detected_at": datetime.utcnow(), + "description": "任务已超过截止日期且未完成。", + "mitigation": "请负责人更新进度、明确阻塞项并给出新的完成时间。", + "evidence": serialize_model(task), + } + ) + return payloads + + def _delayed_project_payloads(self) -> list[dict[str, Any]]: + stmt = select(Project).where( + Project.due_date.is_not(None), + Project.due_date < date.today(), + Project.status.notin_(["验收", "已完成", "复盘", "归档", "关闭", "closed"]), + ) + payloads = [] + for project in self.db.execute(stmt).scalars(): + level = "high" if project.progress_percent < 80 else "medium" + payloads.append( + { + "code": f"RISK-PROJECT-DELAY-{project.id}", + "title": f"项目延期:{project.name}", + "risk_type": "delayed_project", + "risk_level": level, + "status": "open", + "source_domain": "projects", + "source_record_id": str(project.id), + "project_code": project.code, + "owner": project.owner, + "due_date": project.due_date, + "detected_at": datetime.utcnow(), + "description": "项目已超过计划截止日期且未进入完成状态。", + "mitigation": "请项目负责人提交延期原因、资源需求和纠偏计划。", + "evidence": serialize_model(project), + } + ) + return payloads + + def _over_budget_project_payloads(self) -> list[dict[str, Any]]: + stmt = select(Project).where( + Project.budget_amount > 0, + Project.actual_amount > Project.budget_amount, + ) + payloads = [] + for project in self.db.execute(stmt).scalars(): + payloads.append( + { + "code": f"RISK-PROJECT-BUDGET-{project.id}", + "title": f"项目超预算:{project.name}", + "risk_type": "over_budget_project", + "risk_level": "high", + "status": "open", + "source_domain": "projects", + "source_record_id": str(project.id), + "project_code": project.code, + "owner": project.owner, + "due_date": project.due_date, + "detected_at": datetime.utcnow(), + "description": "项目实际成本已超过预算。", + "mitigation": "请复核预算科目、冻结非必要采购并补充审批依据。", + "evidence": serialize_model(project), + } + ) + return payloads + + def _fund_risk_payloads(self) -> list[dict[str, Any]]: + stmt = select(FundAccount).where(FundAccount.current_balance < FundAccount.safety_line) + payloads = [] + for account in self.db.execute(stmt).scalars(): + payloads.append( + { + "code": f"RISK-FUND-{account.id}", + "title": f"资金低于安全线:{account.name}", + "risk_type": "fund_safety_line", + "risk_level": "high", + "status": "open", + "source_domain": "fund-accounts", + "source_record_id": str(account.id), + "owner": None, + "detected_at": datetime.utcnow(), + "description": "账户当前余额低于设置的安全线。", + "mitigation": "请财务确认收付款计划,并优先处理关键项目资金安排。", + "evidence": serialize_model(account), + } + ) + return payloads + + def _supplier_risk_payloads(self) -> list[dict[str, Any]]: + stmt = select(Supplier).where( + (Supplier.blacklist_status != "normal") + | Supplier.risk_level.in_(["medium", "high"]) + ) + payloads = [] + for supplier in self.db.execute(stmt).scalars(): + level = "high" if supplier.blacklist_status != "normal" else supplier.risk_level + payloads.append( + { + "code": f"RISK-SUPPLIER-{supplier.id}", + "title": f"供应商风险:{supplier.name}", + "risk_type": "supplier_risk", + "risk_level": level, + "status": "open", + "source_domain": "suppliers", + "source_record_id": str(supplier.id), + "owner": supplier.contact, + "detected_at": datetime.utcnow(), + "description": "供应商风险等级或黑名单状态需要关注。", + "mitigation": "请采购负责人复核供应商准入、履约和替代方案。", + "evidence": serialize_model(supplier), + } + ) + return payloads diff --git a/app/tools/init_db.py b/app/tools/init_db.py index 252787c..8e855a8 100644 --- a/app/tools/init_db.py +++ b/app/tools/init_db.py @@ -2,14 +2,18 @@ from app.core.database import Base, engine from app.modules.approvals.models import ApprovalRequest from app.modules.audit.models import AuditLog from app.modules.business.models import ( + AttendanceRecord, Expense, FundAccount, + LegacySyncRun, PerformanceMetric, Policy, Procurement, Project, + RiskEvent, Standard, Supplier, + WorkReport, WorkTask, ) @@ -25,6 +29,10 @@ _MODELS = [ Standard, PerformanceMetric, Supplier, + AttendanceRecord, + WorkReport, + RiskEvent, + LegacySyncRun, ] diff --git a/company_ai_local.db b/company_ai_local.db deleted file mode 100644 index 3eb727a82167f88cdde5ac49f09508f1d4436426..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 303104 zcmeI5Yiu0XmFK%jkxf#fT8jvUEbsjj+B zmD%0ZuBw*A5ws2?Tb6@hkogc8Gs(`tGYba+#sLQKcsAo8i^XE+!+Z#KKI{SuWP^Q- zBWo7e%zoG)y8{gN-dnG(>MoLEB-`Zv1xc*F=iYP9{oQlVeH1A@_v|UdAwt8nTB;)y zQwLJ%bm|L&kV>Ua(?7pM|MF}Gw$A^AA)EqjH{%(3P^~=;hB&pqg`RKt> z-gx=J(cICgbjxV#OHOjOIOMfLKc{K^BFTx}}7KGLVzGVE2QNiLG6uRfQk zKCBzH5tiq_eE(=poJx25?AVU#bZy_p1Bo`=dQtsuS^g{cjpoWz>8pLZ^)69#)d{%Y zpQzoh8P)HX=3lxulY5L}8t*efE<40(t4+nw{XX26Xr-re)KHHgf9_K}g0A$L(o8+* z$-SP9Gm8t)o*5KKta8)_E6Kll@}5j?YASu>oTJv8fjZyyPkCyW=bkt<-xGur1AQ@%&AjHgo$x(=J~m$6JMNLnyHlX5g`j*!52#9a#ZTZq%RgFInNY& z!|WHwRdy`Hvlb?0gkFT(R*qUW@9K&o2%kyWKIT5o&xj&P8{BU}G(w}UhP9+SIurRxUp|F13SSZOWlcTxf;q=-7do++tD>Su#15>b;?^7w@F9QG3xn`@Egbhfhu;;ePT9OzF- zPeyZUZ*w-;!%zuC@eF5?m%GGvl#aUAG*vyw6Ey2}%(hMX1hJGF*B6fE;Ox*p`(GK) zJik6NoXI^&$LVk~71sk|^?9;qOPxH$jH(5BPJ%^7`8 z3~%m)mu=XZ-amg`VgA}z(wW?`W9jRro6qAF`&tfpcD?K#ufIW(4`MRNr}0GY&G+6+ z97*IsCa$$qn~AXBVIw~EYLe~q9=;7`xezaDyUlAVSRu)pvN1V8`;KmY_l00ck) z1V8`;KmY_lV7C&OU>(Exf42@V`UL_Y00JNY0w4eaAOHd&00JNY0-ONO{|E#S009sH z0T2KI5C8!X009sH0T9^z1aSV}{bP(Cf&d7B00@8p2!H?xfB*=900@8p&i{x35C8!X z009sH0T2KI5C8!X009u#{RD9S-~D5Z9)bV}fB*=900@8p2!H?xfB*=90M7r20T2KI z5C8!X009sH0T2KI5C8!X*!=`>{@?v$j2?mj2!H?xfB*=900@8p2!H?xfB??_hyf4) z0T2KI5C8!X009sH0T2KI5ZL_$Mraq)*?(jm{T=${txsH!Zhz&plIo?BSZFkeVPELlq^&EOYUxe$(v@k8XvDZk z6hogrE=-%6QX*=B=(R?@K}xzRshV1;YMNRjM3S_6u`+!`=xL-lYaPN`Y8O<8DWOuP zn-*P}*35RpSmBSC44o*pW)ad>TBgpLI98WjrPVJGRVNm&mB>X**1C>ra30)*pQD)}Ork;X7|nUu9}YleDPNrqXjkd)>Im8gd0aYMZLH zsDpFjhlb<-Jk#P4=JYe?W2NQ`oWFsw5Y%fzWSq|edmw= z?)~>a`1)IZ*$3~xN6T)#^@E$={=Q%K;d`%9KW~2j%~fs`0m?3yYZK- z0qx6+SFbX3oKb-~ud9yAnbb{vjXzMXK8Gvkr#ZCq)5p(8pq!sRaz0=}rRUDEs_F`_ z$awGy?gMSQLErm3Kd|?#u&?{9a8Unw`9J{Yr?0;F;?-$3qJBtiJuYc{sSZ&bqs3!V ztjdyDpz&W53lfV1od$_RSa3GfrcG!`i`cfh5`=$Cxhxh6^@>(+&^DEVrfPMup~}Uo zsugNQi4dlyWv>KrZ_uHzRL!Ath>6(b@`Bl&QCibbRJATib+J+{Rcp06X|Rfh9;i?| z9dG8it$PKmDjjRJ?I^2+jsCP)FP7vo5jB~pVzt!ZsJa~*B^z|yl?NSXbZIA;I=icB zl)>~>??Ea^m7|5qQArXCrQ@aY@ltKJSQJI6cs)BjnL2o|N;A&aHSwW2H=DSBH_i5n zLc_96C<(Kzkc)IG^88|J7f6f6@3b^qV7=*??!15?L)5 z%A{6lNDY~TFjq;NsbU|J##$CjWhQr(`{~JO&3Z#U>u0=+BGrSn}M% zx=S6q(mVd%VDd1yGHrHf=JS&5=9!v#)@?JAUX!*Wi=wQRXo?qWa=loWRFy_*sakDR zRGG39c^t$7h~|203!UWQ;7$bIiH#-&9hQPB*Tqt)ESKntLK>3B$2Dm@s?{2oCxxgl z-^OfbOO8Ll9jU{5>bve24_G=_9_0CtrajBBS82wfS(oNqmCl5Yn;!VulsF^g*;272 z$+a8i5Y}L7$p@n2yL zE|zCYvQ&|a*Bfbn4gSVGcYF;l)YWpaPB$3kYOPwU({-g#CG~1qFIGgcMtB-`HeZ7! zx&ll5oPaf$F5`E44VFuiP6~2~ZUH1&t*dljAnO&*d+^+&;~|#0eDE3^%J0q^+%fI= z@hy@v`ANh()5K3I>>Q-jo2EvWX?mb(s_m67-LLag%rmFv78aG|`LoN>vkYI<&$s=P z#O*BY#nNn<9`zUYu=oGt+2>Q)Z)9&~|8@4wte*W(*+0)dpZ({bxSIet1Ogxc0w4ea zAOHd&00JNY0w4eapD+TG!-vxB6$O9&JbrMvnjZM#Xe_A5F6uKA!w1u^LpLFAx9$5C8!X z009sH0T2KI5C8!X*ntEd9-d40e$!|9xLBmWXDhXWRFKC9i>29$ToP;5vACFj%V@bW zTb2rCNy-csmuYKBl*^-u;!3Snl1B!LgYQS(Gg$0@bKdYkaqyiqLsP?Z!3m00ck)1V8`;KmY_l00cl_ zD+C@LelER*A1@zgFY4W&2OZmDlgt+DjfSd6w%BA(*kpK%O@_vY501MXaL@nmrLzAg z`@gccvj3HS_y1pKf0TVM`{S(uLmdzR0T2KI5C8!X009sH0T2KI5CDO@M<6qNFa3#v zOXJ?xG4Ct?34#{Ujc#kyD;n{>?(x13dtZk}hwr675pc;p|9>l${onNT|DUseoBem$ zA7tOjzLovf-SZxrK>!3m00ck)1V8`;KmY_l00ck)1So;M!w;psf0TE;=XtE>In(nz z+Vecp^Sr0$dAR3!Xkz%GaZkT{{{JKP{r~?X`(LvEA^SJu+ zQN6q59>y4gys1!-#?lYr_$X%JGP@bUE8pZgS#(kp$YG*b_Ha<6CO%;LhcX9fils~ok#O7gFsyeE^J zno8d|=cx5&pw4&wQ=Z!8xhGD|_XOd@xS`W%U7mk(eo0t-dRbUJbL!L)VPc${d46u` z#24q5W-6t8M94x{@P$&j9F@8;>5D~4&NGGHF#E-El^x44EyGz0)E1?vR@##li`cf( zA(lqk&OoojxJ=Yj@9HbWQPh^%rAijhoSt7=I58uZjtDV>s^)a*@NZs`_RMm4I0`4} zieo+$R3?T}XBZyZcd1O;_`rkX|QB+wPr;f9tM zPR~cisoPaC&{xpKPH{B@B|VFuQ@F zuzuWFD9J06qq*YY^x6P>G>}XyG_`*NQ@m>23NOr?`;vPU#zAOiasxAA{kW}AlCSSg z946`!4dl?)z71@-b>nv2()_i36PeuM!|B&Fp0z@a;nQAEx0|`ba$zPhcX42m$r70r zZWspg_=NKu_7fGGYmw%3wzCA{taZX1=ub#bMssR!b2iz-Pzglw3}=y-yTo>sj=I(~ zRXxZPH0yQDwoUp3v6LFu7mnrN?9e~^Um4HjP8>_G9}J?nL(}d>ngu=IY{jLc#>*%A zA|PJFYsmv4E)!;WISi$TZz7kN}wabr9hrv9~*H!dL-E+ zi)f~$5AtYqb38%g1Xlt%b!E?{uF$SF4NYx3eSIF;+?Aklf-8YMzdkaY$vsKO>2NX? z*8^hpd9r6qp3rGXtYk;HIRDDfrqjpG8GTL+Z|;PbZP=RLKYv|e{@Pd4ncT5s>FcJO z&*K&QS`K-3z3d*Zzd@1@Vlv35@kH*;_ufn#N#sE&1seP`zS$B$yWxgaq-8|8+fzDZ zz2*n(gX2k??&$2EJz{(<(1~W)bnR6%v(;j|V1F|iW{ttgezUD0Q>;g~ip<;oHG-Z0 zA4ru``_E7Qa`Medao=C>dwTEx9RIWN&y4+xv4=DNDsyb~-|TsJ&;0PeA6^;y<DTY)anb0u>84cE*cRN5d=Mk?5D{$<3%|l)#q%P%oW~sl6t zt<=?~+SW*TYN+&`8p!1iy%nG!X><|eB0Gz1sX86#&Fb1}gE(tS(`av~atk|+w&{>N z7CHSPa_(sKpv#!h{fkY5Rww?4f;ZCf)ybFh$}Hje1;PwJ=w3WRa~wUf2YVD>t^~=& zO-6qCv3P3p1@Q(^V=-;{)^E`LH?cpe`a^&?Do__ou`d@$D)9)@qe(84ejn_V!QTk31n_nvo zUSwEO4o9zVb?NPPAf)MxHS=yh>~5W1Tn18q(esDM%ZPn)qkYjd*nNe2?)KNUc&?1i zn?Oyxo^Ngwb-(Jm)}Ni4pUdRRkEGZ4`T!Xj*0|hk~94LA#(aGF2=*Di%quMK*iAQ~}Fy4^OwqCg~6Vsg6KTD7EkAycny|w*2 z>;S`~8`|dAe|-k|efHJQWODPgw^IY*QkyLFIImv&pWDz}Z>#=(_Db{X!_)LkK5 z_j)b$K7You8e9$saj@MSxbX!sb#uV>nWZS|RS+P(nS;sfE^9_C9;|_o-9B$2xv`|o zw`E(%yIZ~YjnNKVif6jmOiZk|)r@Gi=mx-`rv%0BwpyySb{DsD%Ax4B<|WfwRUFk` zwH1ENAAX*Whi)ixtFmwn%Y+&CJyrK%y8WORwZZlytm407dUNc_a?^~&L4|vCP=eJMz`Ty{H0`LR^AOHd&00JNY0w4eaAOHd&00O&}0M7rrb!5>m5C8!X z009sH0T2KI5C8!X009sP3E=!6vcMAvfB*=900@8p2!H?xfB*=900`_>0yzKg){#ZO zKmY_l00ck)1V8`;KmY_l00clFB!KgO$O2Cw00JNY0w4eaAOHd&00JNY0wAzk3E=#{ zTSpfC0s#;J0T2KI5C8!X009sH0T2LzkO2Gt|569S&v*g>5C8!X009sH0T2KI5C8!X z009u#eFX6SfA@_kdISO>00JNY0w4eaAOHd&00JNY0&xO3|Hp~oGYEhH2!H?xfB*=9 z00@8p2!H?x>^=fG|L?vrMUOxL1V8`;KmY_l00ck)1V8`;Kp;*4=l?hndWJ_0!Z@4hibk3awfKmY_l00ck)1V8`;KmY_l zAWi`1|2PqR1_2NN0T2KI5C8!X009sH0T2Lz-A4fD|J^sH=n)8j00@8p2!H?xfB*=9 z00@8p2*e5C{2wQR&maH-AOHd&00JNY0w4eaAOHd&u=@z${J;Cg6g>g~5C8!X009sH z0T2KI5C8!X0D(9G{QrM(BKQmfAOHd&00JNY0w4eaAOHd&00O&@0N(%azA;6QKmY_l z00ck)1V8`;KmY_l00clFPGEAJq#j89ZfflRj$IwwJNnY-{L1jdLq8f? z++Ub{dvbc;o2lR3_xRrbxYwBY`-wN=I=1)u+Owa}rG;JJDsLMEIW8O zabjtHZh2ljwP`qOimjOzp&p!`Us^aZ zQ$8X{QN=oG8W)MRdBu)pXk=4$mSL|dO>&Vm1CNR^i}j|etu_tYQEW$bx^`kk%_6kd zx}rM5NeX&-;q-ijShu5ZD1T;Y;q=_n3&KcY7(OmIRdad7_+x4K? zBnz++RrR5JbPxIhGt)Pcw|QV!HV+ds-rpaP`cZ1 zz%d*;DSUj7CHv1S#@y#+`LC^2_zZh(&YfYtvOw!`Ppdb>eED`~nCpo@|N5qwZ&q>E zI^j(8=b6fD9$2`5m7LUoO*uoeBq*8#jpU%}M`bEk3Tkz@NJ`X8a;y#s!i8Qi(?8RQhEn z(QVVkLDS4`+p+nuMoP*Ffe@(_L_?(13yCX*Ozs4=cyO@AaIK3pF2w^oB2~gxEC?eq zVI~yjjbm{qIMrT4_g6&!sod`-@ZK`c8-0X;zzTFPF+@Zlu_Y@Yf)X#-uAa!`o}_z|!~S>&6_ognDUduFAM-%c*M_gazL&yG zTxEXsiA?S>YJR-md>E90`rl4A9M+DR30DlF=PPbZTt__HRcE6ng2vp)7)G&K~{>yP@lI>c(2R*MCo(jpGsEW0Ub zK=~7%y}lJHK?UpH05x9RtfwD$T?w+(fV#ZH1Ip*(T^~?~w@viH#gX|<8oU%o5WIu* zxC_q!$rrZbg|}=avGDdD*G@l{$rX#~>+YU%fUyBZpZDy$`^*6q?ox6Kwv>ux(nRDh z*#4H2#Z$Mb221KCVys+nf_ydLET=lEVw+uyCP81Bze!cQju|_>$$jZjGusWrqFd8J zHMg?Wz2o~-y#L?c?FCwa00@8p2!H?xfB*=900@8p2!Oy&B7pP%P8u!v3IZSi0w4ea zAOHd&00JNY0w4ea+arK~|8INT&?8vC_y2a%Xu($y009sH0T2KI5C8!X009sH0T9?80lfd;9yhcC z0T2KI5C8!X009sH0T2KI5CDOlL;&aioitkT6$C&41V8`;KmY_l00ck)1V8`;wnqTx z|Lt)@D-Zwy5C8!X009sH0T2KI5C8!X*hvKN^Zz?(wBRcUfB*=900@8p2!H?xfB*=9 z00?Z40N(#^j~iNn00@8p2!H?xfB*=900@8p2!Oy&B7pP%P8u!v3IZSi0w4eaAOHd& z00JNY0w4ea+arMU|Ms|{6$pR;2!H?xfB*=900@8p2!H?x>?8s>|L>&Hg0CO|0w4ea zAOHd&00JNY0w4eaAh10GIR9^t8(M(?2!H?xfB*=900@8p2!H?xfWS^7!2bRJJ*l@+ z``#RTWaM7@@Jai1!`!n!_1wWjsnosKAA6(w(45&-%0;Or>xH6T)U{%aVcOJ`5>X38uQlooQqomP)znH= z)6^OvlBCs(mFXiwPb0-y>k!sbyP!Hu36(P4wCKvTX0{u~3V*z0=tQwKi;%X`GIiF( zvAX0ct$u;1I&8XakSpj>+f=PprR6#s>O3mkeE)|ZeDj8@`{p}u-2AKW z{``l3NHLvn9~Iov4{ltiMFn2))gS%rJAeFl@4x@S*WdEXK6w8q;BeK{DE@y zIb1nE&7qy2K7Kv|<^1%K^8phoJ$H^(RabaL#)DUIA86AJ`rhC9fxTyiecfk;gZj_Q z2Ld=hef7l`uTHZO^+RgwaY^G#b%^2^Egq9%RhGm8jsKchkXRh(G)Nr6g0rDEZ9+?0 z#J1IyApBd(WwB7GSG0PAwy6{}RjZ2)RW4RltxzjUgfKNNdnJf_gARqIY7U)4OvEOa z7tHRA(wc^%s&z@KiK=~%06M_DCo^ryvou_Tv? zsL4bXtEC1<)$PzI*`VXDJm@&1OFO~T*)GMS)WL&QnsL6ai4V=W*~I<3X|`7s8kTKBNtkVgT%=Qx=NDVMKw2z*r={5f z>rKye=e4}dQ!ePTTA@j%Mh8rk8kK65$ZEMzCbddKYRDXfxk}nh75k7h*0NYCGr6nW zPftc`)*I?sKjU2#xz3H@QbpeIP_qNalII@QUFz7C-tqSalZU~TX|qE!pO<7e&(zGb zZkv(xnzR*J6lJYMQ@mJ{>&3dHsx(qd)oP=n%9NeR;~*A5G}l{O=p+vZcOvjkY&0q8 zuoP6eE|yAVxkOhK(vUPhu1Vujt=70aDMWqwHfB3pa{K}ANFCNw-*vxuz|z6;AkTj^ z?OBGsN;3}4x-{pibS8A%^uX7q#2F#amWm}wuH7()um)od#v0taH%!O_#j;f4dHJ^1 z;DS^XWwr*F<*Hm)OLe)@aMQ_-S%ZU){|akxu{>LnrHWj<-bnjv@Hg(c<7;rCu9k~+ zy1^({Yt>qvt}BHqsaMN-u_B5!!qd33`5G+I6s zTB%Z2>vd66`4M4foZfpnzP+78`M6l5zh^78f>e;NXUF`V!`(Q&*9#4)P*Uk&NmaQ~ zFVMRaU975v)a&(%%&$f^-#O4OOT}7=4H9+^cYfzUCq=zdVv~YGR;!{atF?xv^1&In zbD-lP%C#aJkAXXf9{JtbIo#pXyE`MK*@|2eYt`33o$=S;Yp=Yq_jYg4nQwRN2A%ht zUO3Gy2p!A3Nbk-q!LWt4=?Lzfde4%|Va1w_x$>=5mGZzzbzb`r1pQlk^4G!Fy5% zPhb1eZOl&9m{sjnmA(Vejc)78G@U`M4u5HMK{Xjje$=3El<0jeI~({jj9o0ZO_N_E zs%`rczsCP0PeIxhOTENzY&Y&>;@y6^Omm?m%H`|(M}pmfddKe-+mM_}zYy-fio2<5now2;VXf m+#M9@coamQng;F;dgR;4slWQ&KAr8ieYMXfM1`Ig dict[str, Any]: + return self._data + + +class DummyClient: + calls: list[dict[str, Any]] = [] + + def __init__(self, timeout: int): + self.timeout = timeout + + def __enter__(self) -> "DummyClient": + return self + + def __exit__(self, exc_type, exc, traceback) -> None: + return None + + def post( + self, + url: str, + json: dict[str, Any], + headers: dict[str, str], + ) -> DummyResponse: + self.calls.append({"url": url, "json": json, "headers": headers}) + if url.endswith("/api/v1/agent/ask"): + assert json["context"]["hermes_memory"] == "remembered project preference" + return DummyResponse({"answer": "openclaw final answer", "tool_calls": 1}) + mode = json["context"]["mode"] + if mode == "memory_recall": + return DummyResponse({"answer": "remembered project preference"}) + if mode == "memory_write": + return DummyResponse({"answer": "stored"}) + raise AssertionError(f"unexpected Hermes mode: {mode}") + + +def test_openclaw_hermes_adapter_runs_recall_answer_and_remember(monkeypatch) -> None: + DummyClient.calls = [] + monkeypatch.setattr(adapters.httpx, "Client", DummyClient) + settings = Settings( + model_provider="openclaw_hermes", + openclaw_base_url="http://openclaw.local", + openclaw_api_key="openclaw-key", + hermes_base_url="http://hermes.local", + hermes_api_key="hermes-key", + ) + + result = adapters.OpenClawHermesAdapter(settings).ask( + "summarize project risk", + {"project_code": "P-001"}, + ) + + assert result["answer"] == "openclaw final answer" + assert result["raw"]["pipeline"] == "hermes_recall -> openclaw_answer -> hermes_remember" + assert [call["url"] for call in DummyClient.calls] == [ + "http://hermes.local/api/v1/ask", + "http://openclaw.local/api/v1/agent/ask", + "http://hermes.local/api/v1/ask", + ] + assert DummyClient.calls[0]["headers"]["Authorization"] == "Bearer hermes-key" + assert DummyClient.calls[1]["headers"]["Authorization"] == "Bearer openclaw-key" + assert DummyClient.calls[1]["json"]["context"]["agent_pipeline"] == "openclaw_hermes" diff --git a/tests/test_smoke.py b/tests/test_smoke.py index f4a23b4..221c15a 100644 --- a/tests/test_smoke.py +++ b/tests/test_smoke.py @@ -1,6 +1,7 @@ import json import os import tempfile +from datetime import date, timedelta from pathlib import Path _db = tempfile.NamedTemporaryFile(delete=False, suffix=".db") @@ -154,6 +155,71 @@ def test_approval_gate_for_high_risk_update() -> None: assert update_response.json()["data"]["current_balance"] == 100.0 +def test_new_ledgers_reports_and_risk_events() -> None: + domains_response = client.get("/api/v1/business/domains", headers=headers) + assert domains_response.status_code == 200 + domains = domains_response.json()["domains"] + assert "attendance-records" in domains + assert "work-reports" in domains + assert "risk-events" in domains + + today = date.today() + attendance_response = client.post( + "/api/v1/business/attendance-records", + headers=headers, + json={ + "actor": "pytest", + "data": { + "code": "ATT-SMOKE-001", + "employee_name": "Tester", + "department": "QA", + "work_date": today.isoformat(), + "status": "正常", + }, + }, + ) + assert attendance_response.status_code == 200 + + task_response = client.post( + "/api/v1/business/tasks", + headers=headers, + json={ + "actor": "pytest", + "data": { + "code": "TASK-RISK-001", + "title": "Overdue smoke task", + "owner": "tester", + "status": "待办", + "due_date": (today - timedelta(days=1)).isoformat(), + }, + }, + ) + assert task_response.status_code == 200 + + attendance_summary = client.get("/api/v1/reports/attendance-summary", headers=headers) + assert attendance_summary.status_code == 200 + assert attendance_summary.json()["total"] >= 1 + + report_response = client.post( + "/api/v1/reports/work-reports/generate", + headers=headers, + json={"report_type": "daily", "reporter": "pytest", "actor": "pytest"}, + ) + assert report_response.status_code == 200 + assert report_response.json()["data"]["report_type"] == "daily" + + risk_response = client.post( + "/api/v1/risks/events/generate?actor=pytest", + headers=headers, + ) + assert risk_response.status_code == 200 + assert risk_response.json()["created"] >= 1 + + events_response = client.get("/api/v1/risks/events?status=open", headers=headers) + assert events_response.status_code == 200 + assert any(item["risk_type"] == "overdue_task" for item in events_response.json()["items"]) + + def test_ai_noop_provider() -> None: response = client.post( "/api/v1/ai/ask",