AgentEvalTool/backend/agenteval/storage/db.py
sinohqb 1aa453ef0a feat(intelligent-eval): add cron pool data model and task queue API
Implement Ticket 01 of intelligent eval cron pool architecture (ADR-0007):

- Add 4 new tables: task_queue, cron_pool, config_snapshots, decision_logs
- Implement task enqueueing logic with priority calculation
- Implement task assignment and completion APIs
- Add unit tests (9) and integration tests (7)
- Update CONTEXT.md with new vocabulary
- Add ADR-0007 documenting cron pool architecture decision

All 760 tests passing.
2026-08-12 02:13:21 +08:00

737 lines
24 KiB
Python

"""Database models and session management using SQLModel."""
import json
import uuid
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Optional
import sqlalchemy as sa
from sqlalchemy.pool import StaticPool
from sqlmodel import Field, Relationship, Session, SQLModel, create_engine
DATA_DIR = Path(__file__).resolve().parent.parent.parent.parent / "data"
DATA_DIR.mkdir(parents=True, exist_ok=True)
DATABASE_URL = f"sqlite:///{DATA_DIR / 'agenteval.db'}"
FILES_DIR = DATA_DIR / "files"
FILES_DIR.mkdir(parents=True, exist_ok=True)
engine = create_engine(
DATABASE_URL,
echo=False,
connect_args={"check_same_thread": False},
poolclass=StaticPool,
)
def utc_now() -> datetime:
return datetime.now(timezone.utc)
def as_utc(dt: datetime) -> datetime:
"""Attach UTC tzinfo to a naive datetime.
SQLite round-trips drop tzinfo; stored times are always UTC, so a naive
value read back is restored as UTC before any comparison with utc_now().
"""
return dt if dt.tzinfo is not None else dt.replace(tzinfo=timezone.utc)
def iso_utc(dt: datetime | None) -> str | None:
"""Serialize a datetime to ISO 8601 with UTC timezone suffix.
Guarantees the output always ends with 'Z' or '+00:00' so JavaScript's
Date.parse() interprets it correctly as UTC (no 8-hour local-time offset).
"""
if dt is None:
return None
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt.isoformat(timespec="seconds").replace("+00:00", "Z")
def new_uuid() -> str:
return str(uuid.uuid4())
def _json_dumps(value: Any) -> str:
"""Serialize a JSON column value. ensure_ascii=False keeps CJK readable
in the stored text — the single serialization口径 for all JSON columns."""
return json.dumps(value, ensure_ascii=False)
def _json_loads(raw: str) -> Any:
return json.loads(raw)
class EvalTargetDB(SQLModel, table=True):
"""Database table for evaluation targets."""
__tablename__ = "eval_targets"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
name: str
description: str = ""
platform: str = "ai_digital_employee"
channel_type: str = "tutu-api"
channel_config: str = "{}"
status: str = "active"
created_at: Optional[datetime] = Field(default_factory=utc_now)
updated_at: Optional[datetime] = Field(default_factory=utc_now)
runs: list["EvalRunDB"] = Relationship(
back_populates="target",
sa_relationship_kwargs={"cascade": "all, delete-orphan"},
)
def get_config(self) -> dict[str, Any]:
return _json_loads(self.channel_config)
def set_config(self, config: dict[str, Any]) -> None:
self.channel_config = _json_dumps(config)
class ScenarioDB(SQLModel, table=True):
"""Database table for evaluation scenarios."""
__tablename__ = "scenarios"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
name: str
description: str = ""
tags: str = "[]"
cases: str = "[]"
llm_config: Optional[str] = None
version: int = Field(default=1)
created_at: Optional[datetime] = Field(default_factory=utc_now)
updated_at: Optional[datetime] = Field(default_factory=utc_now)
runs: list["EvalRunDB"] = Relationship(
back_populates="scenario",
sa_relationship_kwargs={"cascade": "all, delete-orphan"},
)
def get_tags(self) -> list[str]:
return _json_loads(self.tags)
def set_tags(self, tags: list[str]) -> None:
self.tags = _json_dumps(tags)
def get_cases(self) -> list[dict[str, Any]]:
return _json_loads(self.cases)
def set_cases(self, cases: list[dict[str, Any]]) -> None:
self.cases = _json_dumps(cases)
def get_llm_config(self) -> Optional[dict[str, Any]]:
return _json_loads(self.llm_config) if self.llm_config else None
def set_llm_config(self, config: Optional[dict[str, Any]]) -> None:
self.llm_config = _json_dumps(config) if config else None
class ModelConfigDB(SQLModel, table=True):
"""Reusable external model connection configuration."""
__tablename__ = "model_configs"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
name: str = Field(index=True, unique=True)
provider: str = "openai_compatible"
capability: str = Field(index=True)
endpoint_url: str
model_name: Optional[str] = None
vendor_name: str = ""
input_modalities: str = '["text"]'
output_modalities: str = '["text"]'
context_window: Optional[int] = None
max_output_tokens: Optional[int] = None
supports_streaming: bool = False
supports_tool_calling: bool = False
supports_structured_output: bool = False
supports_reasoning: bool = False
region: str = ""
documentation_url: Optional[str] = None
api_key_encrypted: Optional[str] = None
enabled: bool = True
is_default: bool = False
is_analysis_default: bool = False
description: str = ""
created_at: Optional[datetime] = Field(default_factory=utc_now)
updated_at: Optional[datetime] = Field(default_factory=utc_now)
def get_input_modalities(self) -> list[str]:
return _json_loads(self.input_modalities)
def get_output_modalities(self) -> list[str]:
return _json_loads(self.output_modalities)
def set_modalities(self, input_modalities: list[str], output_modalities: list[str]) -> None:
self.input_modalities = _json_dumps(input_modalities)
self.output_modalities = _json_dumps(output_modalities)
class ScenarioModelBindingDB(SQLModel, table=True):
"""Bind one model configuration to a purpose within a scenario."""
__tablename__ = "scenario_model_bindings"
scenario_id: str = Field(foreign_key="scenarios.id", primary_key=True)
purpose: str = Field(primary_key=True)
model_config_id: str = Field(foreign_key="model_configs.id", index=True)
class CampaignDB(SQLModel, table=True):
"""Database table for evaluation campaigns (评估活动)."""
__tablename__ = "campaigns"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
name: str
target_id: Optional[str] = Field(default=None, foreign_key="eval_targets.id")
window_seconds: int
time_scale: float = 1.0
plan: str = "[]"
status: str = "planned"
started_at: Optional[datetime] = None
completed_at: Optional[datetime] = None
created_at: Optional[datetime] = Field(default_factory=utc_now)
summary: Optional[str] = None
analysis_model_config_id: Optional[str] = None
exploration_seeds: Optional[str] = None
exploration_budget: Optional[str] = None
last_patrolled_at: Optional[datetime] = None
def get_plan(self) -> list[dict[str, Any]]:
return _json_loads(self.plan)
def set_plan(self, plan: list[dict[str, Any]]) -> None:
self.plan = _json_dumps(plan)
def get_summary(self) -> Optional[dict[str, Any]]:
return _json_loads(self.summary) if self.summary else None
def set_summary(self, summary: dict[str, Any]) -> None:
self.summary = _json_dumps(summary)
def get_exploration_seeds(self) -> Optional[dict[str, Any]]:
return _json_loads(self.exploration_seeds) if self.exploration_seeds else None
def set_exploration_seeds(self, seeds: dict[str, Any]) -> None:
self.exploration_seeds = _json_dumps(seeds)
def get_exploration_budget(self) -> Optional[dict[str, Any]]:
return _json_loads(self.exploration_budget) if self.exploration_budget else None
def set_exploration_budget(self, budget: dict[str, Any]) -> None:
self.exploration_budget = _json_dumps(budget)
class CampaignAnalysisDB(SQLModel, table=True):
"""One row per campaign holding its intelligent analysis (智能分析) state."""
__tablename__ = "campaign_analyses"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
campaign_id: str = Field(unique=True, index=True)
status: str = "generating"
result: Optional[str] = None
model_config_id: Optional[str] = None
error: Optional[str] = None
triggered_by: str = "manual"
recovery_attempts: int = 0
created_at: Optional[datetime] = Field(default_factory=utc_now)
updated_at: Optional[datetime] = Field(default_factory=utc_now)
def get_result(self) -> Optional[dict[str, Any]]:
return _json_loads(self.result) if self.result else None
def set_result(self, result: dict[str, Any]) -> None:
self.result = _json_dumps(result)
class CampaignPeriodComparisonDB(SQLModel, table=True):
"""One row per campaign holding its period-comparison (周期对比) state.
Mirrors ``campaign_analyses`` (one row, upserted per campaign), plus the
baseline campaign this comparison was computed against.
"""
__tablename__ = "campaign_period_comparisons"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
campaign_id: str = Field(unique=True, index=True)
baseline_campaign_id: str = Field(index=True)
status: str = "generating"
result: Optional[str] = None
model_config_id: Optional[str] = None
error: Optional[str] = None
triggered_by: str = "manual"
recovery_attempts: int = 0
created_at: Optional[datetime] = Field(default_factory=utc_now)
updated_at: Optional[datetime] = Field(default_factory=utc_now)
def get_result(self) -> Optional[dict[str, Any]]:
return _json_loads(self.result) if self.result else None
def set_result(self, result: dict[str, Any]) -> None:
self.result = _json_dumps(result)
class ExplorationSessionDB(SQLModel, table=True):
"""One virtual-user exploration session (探索会话) within a campaign.
Independent entity, deliberately NOT an EvalRun: exploration outcomes feed
the 体验判定 evidence line and must not pollute pass-rate semantics
(ADR-0001 comparability / ADR-0002).
"""
__tablename__ = "exploration_sessions"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
campaign_id: str = Field(index=True, foreign_key="campaigns.id")
target_id: str = Field(foreign_key="eval_targets.id")
persona: str = "{}"
goal: str = ""
seed_ref: Optional[str] = None
status: str = "running"
triggered_by: str = "auto"
experience: Optional[str] = None
judge_review: Optional[str] = None
turn_count: int = 0
error: Optional[str] = None
created_at: Optional[datetime] = Field(default_factory=utc_now)
closed_at: Optional[datetime] = None
def get_persona(self) -> dict[str, Any]:
return _json_loads(self.persona)
def set_persona(self, persona: dict[str, Any]) -> None:
self.persona = _json_dumps(persona)
def get_seed_ref(self) -> Optional[dict[str, Any]]:
return _json_loads(self.seed_ref) if self.seed_ref else None
def set_seed_ref(self, seed_ref: dict[str, Any]) -> None:
self.seed_ref = _json_dumps(seed_ref)
def get_experience(self) -> Optional[dict[str, Any]]:
return _json_loads(self.experience) if self.experience else None
def set_experience(self, experience: dict[str, Any]) -> None:
self.experience = _json_dumps(experience)
def get_judge_review(self) -> Optional[dict[str, Any]]:
return _json_loads(self.judge_review) if self.judge_review else None
def set_judge_review(self, judge_review: dict[str, Any]) -> None:
self.judge_review = _json_dumps(judge_review)
class ExplorationMessageDB(SQLModel, table=True):
"""One chat message inside an exploration session (role/content form).
Parallel to ``eval_runs`` turns but never shared with them. A completed
question-answer pair contributes two rows (user + assistant); the reply
row carries ``latency_ms``.
"""
__tablename__ = "exploration_messages"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
session_id: str = Field(index=True, foreign_key="exploration_sessions.id")
round_index: int = 0
role: str = "user"
content: str = ""
latency_ms: Optional[int] = None
created_at: Optional[datetime] = Field(default_factory=utc_now)
class EvalRunDB(SQLModel, table=True):
"""Database table for evaluation runs."""
__tablename__ = "eval_runs"
__table_args__ = (
sa.Index(
"uq_eval_runs_campaign_occurrence",
"campaign_id",
"campaign_plan_index",
"campaign_occurrence_index",
unique=True,
),
)
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
target_id: Optional[str] = Field(default=None, foreign_key="eval_targets.id")
scenario_id: Optional[str] = Field(default=None, foreign_key="scenarios.id")
scenario_version: int = Field(default=1)
campaign_id: Optional[str] = Field(default=None, foreign_key="campaigns.id")
campaign_plan_index: Optional[int] = Field(default=None, ge=0)
campaign_occurrence_index: Optional[int] = Field(default=None, ge=0)
status: str = "pending"
triggered_by: str = Field(default="manual")
started_at: Optional[datetime] = Field(default_factory=utc_now)
completed_at: Optional[datetime] = None
summary: Optional[str] = None
target: Optional[EvalTargetDB] = Relationship(back_populates="runs")
scenario: Optional[ScenarioDB] = Relationship(back_populates="runs")
turns: list["TurnDB"] = Relationship(
back_populates="run",
sa_relationship_kwargs={"cascade": "all, delete-orphan"},
)
results: list["EvalResultDB"] = Relationship(
back_populates="run",
sa_relationship_kwargs={"cascade": "all, delete-orphan"},
)
def get_summary(self) -> Optional[dict[str, Any]]:
return _json_loads(self.summary) if self.summary else None
def set_summary(self, summary: dict[str, Any]) -> None:
self.summary = _json_dumps(summary)
class TurnDB(SQLModel, table=True):
"""Database table for conversation turns."""
__tablename__ = "turns"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
run_id: Optional[str] = Field(default=None, foreign_key="eval_runs.id")
case_id: str
round_index: int
sent_message: str = "{}"
sent_at: Optional[datetime] = Field(default_factory=utc_now)
question_msg_id: Optional[str] = None
reply: Optional[str] = None
received_at: Optional[datetime] = None
latency_ms: Optional[int] = None
run: Optional[EvalRunDB] = Relationship(back_populates="turns")
def get_sent_message(self) -> dict[str, Any]:
return _json_loads(self.sent_message)
def set_sent_message(self, message: dict[str, Any]) -> None:
self.sent_message = _json_dumps(message)
def get_reply(self) -> Optional[dict[str, Any]]:
return _json_loads(self.reply) if self.reply else None
def set_reply(self, reply: Optional[dict[str, Any]]) -> None:
self.reply = _json_dumps(reply) if reply else None
class EvalResultDB(SQLModel, table=True):
"""Database table for evaluation results."""
__tablename__ = "eval_results"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
run_id: Optional[str] = Field(default=None, foreign_key="eval_runs.id")
case_id: str
turn_id: str
rule_type: str
passed: bool
score: Optional[float] = None
reason: str = ""
run: Optional[EvalRunDB] = Relationship(back_populates="results")
class FileCategoryDB(SQLModel, table=True):
"""Database table for file categories (tree structure via parent_id)."""
__tablename__ = "file_categories"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
name: str
parent_id: Optional[str] = Field(
default=None,
foreign_key="file_categories.id",
)
created_at: Optional[datetime] = Field(default_factory=utc_now)
updated_at: Optional[datetime] = Field(default_factory=utc_now)
parent: Optional["FileCategoryDB"] = Relationship(
back_populates="children",
sa_relationship_kwargs={"remote_side": "FileCategoryDB.id"},
)
children: list["FileCategoryDB"] = Relationship(
back_populates="parent",
sa_relationship_kwargs={"cascade": "all, delete-orphan"},
)
files: list["FileRecordDB"] = Relationship(
back_populates="category",
sa_relationship_kwargs={"cascade": "all, delete-orphan"},
)
class FileRecordDB(SQLModel, table=True):
"""Database table for uploaded file records."""
__tablename__ = "file_records"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
original_name: str
storage_name: str
category_id: Optional[str] = Field(
default=None,
foreign_key="file_categories.id",
)
file_size: int = 0
mime_type: str = ""
file_ext: str = ""
created_at: Optional[datetime] = Field(default_factory=utc_now)
category: Optional[FileCategoryDB] = Relationship(back_populates="files")
class IntelligentEvalDB(SQLModel, table=True):
"""Intelligent evaluation (智能评估) — independent entity, peer to Campaign."""
__tablename__ = "intelligent_evals"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
name: str
target_id: str = Field(foreign_key="eval_targets.id")
status: str = "draft"
# User input (四件套)
goal: str = ""
seeds: str = "{}"
intent: str = ""
role_description: str = ""
# Coarse plan (OpenClaw produces)
plan: Optional[str] = None
plan_feedback: Optional[str] = None
# Time window
time_window_hours: int = 24
# Report
report: Optional[str] = None
# Timestamps
created_at: Optional[datetime] = Field(default_factory=utc_now)
updated_at: Optional[datetime] = Field(default_factory=utc_now)
started_at: Optional[datetime] = None
completed_at: Optional[datetime] = None
def get_seeds(self) -> dict[str, Any]:
return _json_loads(self.seeds)
def set_seeds(self, seeds: dict[str, Any]) -> None:
self.seeds = _json_dumps(seeds)
def get_plan(self) -> Optional[dict[str, Any]]:
return _json_loads(self.plan) if self.plan else None
def set_plan(self, plan: dict[str, Any]) -> None:
self.plan = _json_dumps(plan)
def get_report(self) -> Optional[dict[str, Any]]:
return _json_loads(self.report) if self.report else None
def set_report(self, report: dict[str, Any]) -> None:
self.report = _json_dumps(report)
class IntelligentEvalSessionDB(SQLModel, table=True):
"""A virtual-user session within an intelligent evaluation."""
__tablename__ = "intelligent_eval_sessions"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
eval_id: str = Field(index=True, foreign_key="intelligent_evals.id")
target_id: str = Field(foreign_key="eval_targets.id")
persona: str = "{}"
goal: str = ""
dimension: Optional[str] = None
status: str = "running"
verdict: Optional[str] = None
turn_count: int = 0
created_at: Optional[datetime] = Field(default_factory=utc_now)
closed_at: Optional[datetime] = None
def get_persona(self) -> dict[str, Any]:
return _json_loads(self.persona)
def set_persona(self, persona: dict[str, Any]) -> None:
self.persona = _json_dumps(persona)
def get_verdict(self) -> Optional[dict[str, Any]]:
return _json_loads(self.verdict) if self.verdict else None
def set_verdict(self, verdict: dict[str, Any]) -> None:
self.verdict = _json_dumps(verdict)
class IntelligentEvalMessageDB(SQLModel, table=True):
"""One chat message inside an intelligent eval session."""
__tablename__ = "intelligent_eval_messages"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
session_id: str = Field(index=True, foreign_key="intelligent_eval_sessions.id")
role: str = "user"
content: str = ""
latency_ms: Optional[int] = None
created_at: Optional[datetime] = Field(default_factory=utc_now)
class IntelligentEvalTaskQueueDB(SQLModel, table=True):
"""Task queue for intelligent evaluations (任务队列).
Platform scans executing evals every minute and enqueues tasks for
OpenClaw workers to pick up.
"""
__tablename__ = "intelligent_eval_task_queue"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
eval_id: str = Field(index=True, foreign_key="intelligent_evals.id")
# Task status
status: str = Field(index=True) # pending / assigned / completed / failed
priority: int = Field(index=True) # Lower is higher priority
reason: str # Why this task needs attention (e.g., "slot_due", "all_sessions_completed")
# Assignment info
assigned_cron_id: Optional[str] = None
assigned_at: Optional[datetime] = None
# Completion info
completed_at: Optional[datetime] = None
error: Optional[str] = None
# Timestamps
created_at: Optional[datetime] = Field(default_factory=utc_now)
updated_at: Optional[datetime] = Field(default_factory=utc_now)
__table_args__ = (
sa.Index("idx_task_queue_status_priority", "status", "priority"),
sa.Index("idx_task_queue_eval_status", "eval_id", "status"),
)
class OpenClawCronPoolDB(SQLModel, table=True):
"""OpenClaw cron pool state (平台侧记录).
Platform tracks which crons exist and their current state.
"""
__tablename__ = "openclaw_cron_pool"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
openclaw_cron_id: str = Field(unique=True, index=True) # OpenClaw's cron job ID
# State
status: str = Field(index=True) # idle / busy / stuck
current_eval_id: Optional[str] = Field(foreign_key="intelligent_evals.id")
# Heartbeat
last_active_at: datetime
last_task_at: Optional[datetime] = None
# Stats
total_tasks_completed: int = 0
total_tasks_failed: int = 0
# Timestamps
created_at: Optional[datetime] = Field(default_factory=utc_now)
updated_at: Optional[datetime] = Field(default_factory=utc_now)
__table_args__ = (sa.Index("idx_cron_pool_status_last_active", "status", "last_active_at"),)
class IntelligentEvalConfigSnapshotDB(SQLModel, table=True):
"""Config snapshot for intelligent evaluations (配置快照).
Automatically saved when eval is created, plan is submitted, or config is updated.
"""
__tablename__ = "intelligent_eval_config_snapshots"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
eval_id: str = Field(index=True, foreign_key="intelligent_evals.id")
# Snapshot type
snapshot_type: str # created / plan_submitted / config_updated
# Config snapshot
goal: str
seeds: str # JSON
intent: str
role_description: str
time_window_hours: int
plan: Optional[str] = None # JSON, coarse plan snapshot
# Metadata
created_at: Optional[datetime] = Field(default_factory=utc_now)
created_by: str = "user" # user / openclaw
__table_args__ = (sa.Index("idx_config_snapshots_eval_created", "eval_id", "created_at"),)
def get_seeds(self) -> dict[str, Any]:
return _json_loads(self.seeds)
def set_seeds(self, seeds: dict[str, Any]) -> None:
self.seeds = _json_dumps(seeds)
def get_plan(self) -> Optional[dict[str, Any]]:
return _json_loads(self.plan) if self.plan else None
def set_plan(self, plan: dict[str, Any]) -> None:
self.plan = _json_dumps(plan)
class IntelligentEvalDecisionLogDB(SQLModel, table=True):
"""Decision log for intelligent evaluations (决策日志).
Records every decision made by OpenClaw workers (execute_session / wait / start_analysis).
"""
__tablename__ = "intelligent_eval_decision_logs"
id: Optional[str] = Field(default_factory=new_uuid, primary_key=True)
eval_id: str = Field(index=True, foreign_key="intelligent_evals.id")
# Decision info
decision_type: str # execute_session / wait / start_analysis
reason: str # Why this decision was made
context: str # JSON, decision context (current time slot, completed sessions, etc.)
# Execution info
cron_id: str # Which cron made this decision
# Timestamp
created_at: Optional[datetime] = Field(default_factory=utc_now)
__table_args__ = (sa.Index("idx_decision_logs_eval_created", "eval_id", "created_at"),)
def get_context(self) -> dict[str, Any]:
return _json_loads(self.context)
def set_context(self, context: dict[str, Any]) -> None:
self.context = _json_dumps(context)
def init_db() -> None:
SQLModel.metadata.create_all(engine)
def get_session() -> Session:
return Session(engine)
def get_session_context():
"""Context manager that creates and properly closes a database session."""
session = Session(engine)
try:
yield session
finally:
session.close()