Files
2026-08-22 07:25:48 +02:00

1059 lines
41 KiB
Python

"""Self-improvement services — signal collection, pattern detection,
proposal lifecycle, evaluation, activation, impact measurement.
Builds on existing systems:
- ai_proactive: ContextLog, ProactiveSuggestion
- automation: AgentRun, AgentRunStep, AgentVersion
- app.core.approval: create_approval_request
- app.models.audit: AuditLog
- app.models.workflow: WorkflowInstance
- app.ai.llm_client: llm_complete
"""
from __future__ import annotations
import logging
import uuid
from datetime import UTC, datetime, timedelta
from typing import Any
from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession
from app.ai.llm_client import llm_complete
from app.plugins.builtins.self_improvement.models import (
ImprovementPattern,
ImprovementProposal,
ImprovementSignal,
ImpactMeasurement,
PROPOSAL_STATUSES,
)
logger = logging.getLogger(__name__)
# ──────────────────────────────────────────────────────────────────────────
# J-SIGNAL: Signal Collection
# ──────────────────────────────────────────────────────────────────────────
async def collect_signals(
db: AsyncSession,
tenant_id: uuid.UUID,
since: datetime | None = None,
limit: int = 100,
) -> dict[str, Any]:
"""Collect improvement signals from existing system data.
Queries real usage data from AgentRun, WorkflowInstance,
ProactiveSuggestion, AuditLog — stores references/summaries,
not personal data copies.
"""
if since is None:
since = datetime.now(UTC) - timedelta(days=7)
collected: list[ImprovementSignal] = []
# 1. Agent runs with failures or retries
try:
from app.plugins.builtins.automation.models import AgentRun, AgentRunStep
# Failed agent runs
result = await db.execute(
select(AgentRun)
.where(
AgentRun.tenant_id == tenant_id,
AgentRun.created_at >= since,
AgentRun.status.in_(["failed", "error", "timeout"]),
)
.order_by(AgentRun.created_at.desc())
.limit(limit)
)
for run in result.scalars().all():
sig = ImprovementSignal(
tenant_id=tenant_id,
source_type="agent_run",
source_ref_id=run.id,
source_metadata={
"agent_id": str(run.agent_id) if run.agent_id else None,
"status": run.status,
"duration_seconds": getattr(run, "duration_seconds", None),
},
summary=f"Agent run failed with status '{run.status}'",
signal_kind="failure",
severity="error",
confidence=0.8,
)
db.add(sig)
collected.append(sig)
# Agent runs with many steps (potential retry/bottleneck)
result = await db.execute(
select(AgentRunStep.agent_run_id, func.count(AgentRunStep.id).label("step_count"))
.where(
AgentRunStep.tenant_id == tenant_id,
AgentRunStep.created_at >= since,
)
.group_by(AgentRunStep.agent_run_id)
.having(func.count(AgentRunStep.id) >= 5)
.limit(limit)
)
for row in result.all():
sig = ImprovementSignal(
tenant_id=tenant_id,
source_type="agent_run",
source_ref_id=row.agent_run_id,
source_metadata={"step_count": row.step_count},
summary=f"Agent run had {row.step_count} steps (potential bottleneck)",
signal_kind="retry",
severity="warning",
confidence=0.6,
)
db.add(sig)
collected.append(sig)
except Exception:
logger.warning("Failed to collect agent run signals", exc_info=True)
# 2. Workflow instances with failures
try:
from app.models.workflow import WorkflowInstance
result = await db.execute(
select(WorkflowInstance)
.where(
WorkflowInstance.tenant_id == tenant_id,
WorkflowInstance.created_at >= since,
WorkflowInstance.status.in_(["failed", "error", "cancelled"]),
)
.order_by(WorkflowInstance.created_at.desc())
.limit(limit)
)
for wf in result.scalars().all():
sig = ImprovementSignal(
tenant_id=tenant_id,
source_type="workflow_run",
source_ref_id=wf.id,
source_metadata={
"workflow_id": str(wf.workflow_id) if hasattr(wf, "workflow_id") else None,
"status": wf.status,
},
summary=f"Workflow instance failed with status '{wf.status}'",
signal_kind="failure",
severity="error",
confidence=0.7,
)
db.add(sig)
collected.append(sig)
except Exception:
logger.warning("Failed to collect workflow signals", exc_info=True)
# 3. Dismissed proactive suggestions
try:
from app.plugins.builtins.ai_proactive.models import ProactiveSuggestion
result = await db.execute(
select(ProactiveSuggestion)
.where(
ProactiveSuggestion.tenant_id == tenant_id,
ProactiveSuggestion.created_at >= since,
ProactiveSuggestion.is_dismissed.is_(True),
)
.order_by(ProactiveSuggestion.created_at.desc())
.limit(limit)
)
for sug in result.scalars().all():
sig = ImprovementSignal(
tenant_id=tenant_id,
source_type="proactive_suggestion",
source_ref_id=sug.id,
source_metadata={
"suggestion_type": sug.suggestion_type,
"entity_type": sug.entity_type,
},
summary=f"Proactive suggestion '{sug.title}' was dismissed",
signal_kind="dismissal",
severity="info",
confidence=0.5,
)
db.add(sig)
collected.append(sig)
except Exception:
logger.warning("Failed to collect proactive suggestion signals", exc_info=True)
# 4. Audit log entries with corrections (update/delete patterns)
try:
from app.models.audit import AuditLog
result = await db.execute(
select(AuditLog.entity_type, AuditLog.action, func.count(AuditLog.id).label("cnt"))
.where(
AuditLog.tenant_id == tenant_id,
AuditLog.created_at >= since,
AuditLog.action.in_(["update", "delete"]),
)
.group_by(AuditLog.entity_type, AuditLog.action)
.having(func.count(AuditLog.id) >= 3)
.limit(limit)
)
for row in result.all():
sig = ImprovementSignal(
tenant_id=tenant_id,
source_type="audit_log",
source_ref_id=None,
source_metadata={
"entity_type": row.entity_type,
"action": row.action,
"count": row.cnt,
},
summary=f"{row.cnt} '{row.action}' operations on '{row.entity_type}' (potential correction pattern)",
signal_kind="correction",
severity="info",
confidence=0.5,
)
db.add(sig)
collected.append(sig)
except Exception:
logger.warning("Failed to collect audit log signals", exc_info=True)
await db.flush()
return {
"collected": len(collected),
"signals": [
{
"id": str(s.id),
"source_type": s.source_type,
"signal_kind": s.signal_kind,
"severity": s.severity,
"summary": s.summary,
"confidence": s.confidence,
}
for s in collected
],
}
# ──────────────────────────────────────────────────────────────────────────
# J-PATTERN: Pattern Detection
# ──────────────────────────────────────────────────────────────────────────
async def detect_patterns(
db: AsyncSession,
tenant_id: uuid.UUID,
min_occurrences: int = 2,
) -> dict[str, Any]:
"""Detect recurring patterns from collected signals.
Groups signals by source_type + signal_kind + target metadata
to find bottlenecks, repetitive corrections, dismissed suggestions.
"""
# Get unpatterned signals
result = await db.execute(
select(ImprovementSignal)
.where(
ImprovementSignal.tenant_id == tenant_id,
ImprovementSignal.pattern_id.is_(None),
)
.order_by(ImprovementSignal.created_at.desc())
)
signals = result.scalars().all()
if not signals:
return {"patterns_created": 0, "patterns": []}
# Group signals by (source_type, signal_kind, target identifier from metadata)
groups: dict[str, list[ImprovementSignal]] = {}
for sig in signals:
# Create a grouping key from source type + kind + a target identifier
meta = sig.source_metadata or {}
target_key = meta.get("agent_id") or meta.get("workflow_id") or meta.get("entity_type") or meta.get("suggestion_type") or "unknown"
key = f"{sig.source_type}:{sig.signal_kind}:{target_key}"
groups.setdefault(key, []).append(sig)
patterns_created: list[ImprovementPattern] = []
for key, group_signals in groups.items():
if len(group_signals) < min_occurrences:
continue
parts = key.split(":", 2)
source_type = parts[0] if len(parts) > 0 else "unknown"
signal_kind = parts[1] if len(parts) > 1 else "observation"
target_name = parts[2] if len(parts) > 2 else "unknown"
# Determine target_type from source
target_type_map = {
"agent_run": "agent",
"workflow_run": "workflow",
"proactive_suggestion": "trigger",
"audit_log": "agent",
}
target_type = target_type_map.get(source_type, "agent")
# Build evidence refs (references, not data copies)
evidence_refs = [
{
"signal_id": str(s.id),
"source_type": s.source_type,
"source_ref_id": str(s.source_ref_id) if s.source_ref_id else None,
"summary": s.summary,
}
for s in group_signals
]
avg_confidence = sum(s.confidence for s in group_signals) / len(group_signals)
pattern_kind_map = {
"failure": "retry_bottleneck",
"retry": "retry_bottleneck",
"correction": "manual_correction",
"dismissal": "suggestion_dismissal",
}
pattern_kind = pattern_kind_map.get(signal_kind, "repetitive_handoff")
pattern = ImprovementPattern(
tenant_id=tenant_id,
pattern_kind=pattern_kind,
title=f"{pattern_kind.replace('_', ' ').title()}: {target_name} ({len(group_signals)} occurrences)",
description=f"Detected {len(group_signals)} '{signal_kind}' signals for '{target_name}' from '{source_type}'. "
+ "; ".join(s.summary for s in group_signals[:3]),
target_type=target_type,
target_name=target_name,
evidence_refs=evidence_refs,
occurrence_count=len(group_signals),
confidence=avg_confidence,
status="detected",
proposed_action=_suggest_action(pattern_kind, target_type, target_name, len(group_signals)),
)
db.add(pattern)
await db.flush()
# Link signals to pattern
for sig in group_signals:
sig.pattern_id = pattern.id
patterns_created.append(pattern)
await db.flush()
return {
"patterns_created": len(patterns_created),
"patterns": [
{
"id": str(p.id),
"pattern_kind": p.pattern_kind,
"title": p.title,
"target_type": p.target_type,
"target_name": p.target_name,
"occurrence_count": p.occurrence_count,
"confidence": p.confidence,
"proposed_action": p.proposed_action,
}
for p in patterns_created
],
}
def _suggest_action(pattern_kind: str, target_type: str, target_name: str, count: int) -> str:
"""Generate a human-readable proposed action for a pattern."""
if pattern_kind == "retry_bottleneck":
return f"Review {target_type} '{target_name}': {count} retries/failures detected. Consider adjusting system prompt, tools, or timeout settings."
elif pattern_kind == "manual_correction":
return f"Review {target_type} '{target_name}': {count} manual corrections detected. Consider improving output quality or adding validation."
elif pattern_kind == "suggestion_dismissal":
return f"Review trigger '{target_name}': {count} dismissed suggestions. Consider adjusting confidence threshold or suggestion relevance."
else:
return f"Review {target_type} '{target_name}': {count} repetitive handoffs detected. Consider automation or workflow adjustment."
# ──────────────────────────────────────────────────────────────────────────
# J-PROP: Improvement Proposals
# ──────────────────────────────────────────────────────────────────────────
async def create_proposal(
db: AsyncSession,
tenant_id: uuid.UUID,
*,
pattern_id: uuid.UUID | None = None,
title: str,
description: str,
target_type: str,
target_ref_id: uuid.UUID | None = None,
target_name: str | None = None,
proposed_config: dict | None = None,
rationale: str = "",
expected_benefit: str = "",
risk_assessment: str = "",
user_id: uuid.UUID | None = None,
) -> ImprovementProposal:
"""Create a new improvement proposal in draft status."""
# Capture previous config for rollback
previous_config: dict[str, Any] = {}
if target_ref_id and target_type == "agent":
try:
from app.plugins.builtins.automation.models import AgentDefinition
result = await db.execute(
select(AgentDefinition).where(
AgentDefinition.tenant_id == tenant_id,
AgentDefinition.id == target_ref_id,
)
)
agent = result.scalar_one_or_none()
if agent:
previous_config = {
"system_prompt": agent.system_prompt,
"tool_ids": agent.tool_ids,
"temperature": agent.temperature,
"max_tokens": agent.max_tokens,
"max_steps": agent.max_steps,
}
except Exception:
logger.warning("Failed to capture previous agent config", exc_info=True)
# Get evidence from pattern if linked
evidence_refs: list = []
if pattern_id:
pat_result = await db.execute(
select(ImprovementPattern).where(
ImprovementPattern.tenant_id == tenant_id,
ImprovementPattern.id == pattern_id,
)
)
pattern = pat_result.scalar_one_or_none()
if pattern:
evidence_refs = pattern.evidence_refs
# Update pattern status
pattern.status = "proposal_created"
proposal = ImprovementProposal(
tenant_id=tenant_id,
owner_id=user_id,
pattern_id=pattern_id,
title=title,
description=description,
target_type=target_type,
target_ref_id=target_ref_id,
target_name=target_name,
proposed_config=proposed_config or {},
previous_config=previous_config,
evidence_refs=evidence_refs,
rationale=rationale,
expected_benefit=expected_benefit,
risk_assessment=risk_assessment,
status="draft",
)
db.add(proposal)
await db.flush()
return proposal
# ──────────────────────────────────────────────────────────────────────────
# J-EVAL: Evaluation / Dry-Run
# ──────────────────────────────────────────────────────────────────────────
async def evaluate_proposal(
db: AsyncSession,
tenant_id: uuid.UUID,
proposal_id: uuid.UUID,
) -> dict[str, Any]:
"""Evaluate a proposal via LLM-based dry-run assessment.
No external side effects — the LLM assesses the proposed change
against the evidence and expected benefit.
"""
result = await db.execute(
select(ImprovementProposal).where(
ImprovementProposal.tenant_id == tenant_id,
ImprovementProposal.id == proposal_id,
)
)
proposal = result.scalar_one_or_none()
if not proposal:
return {"error": "Proposal not found"}
proposal.status = "evaluating"
await db.flush()
# Build evaluation prompt
eval_prompt = f"""You are an AI improvement evaluator for a CRM system.
Assess the following improvement proposal. Consider:
1. Is the rationale sound?
2. Is the expected benefit realistic?
3. Are there risks not mentioned?
4. Would you recommend approval?
Return JSON:
{{
"score": 0.0-1.0,
"assessment": "...",
"risks_identified": ["..."],
"recommendation": "approve" | "reject" | "needs_review",
"test_scenarios": ["..."]
}}
Proposal:
- Title: {proposal.title}
- Target: {proposal.target_type} ({proposal.target_name or 'N/A'})
- Description: {proposal.description}
- Rationale: {proposal.rationale}
- Expected Benefit: {proposal.expected_benefit}
- Risk Assessment: {proposal.risk_assessment}
- Evidence Count: {len(proposal.evidence_refs)}
- Proposed Config: {proposal.proposed_config}
"""
try:
response = await llm_complete(
model="openai/gpt-4o-mini",
messages=[
{"role": "system", "content": "You are an improvement evaluator. Return only valid JSON."},
{"role": "user", "content": eval_prompt},
],
temperature=0.2,
max_tokens=1000,
tenant_id=tenant_id,
db=db,
)
import json
eval_result = json.loads(response.get("content", "{}"))
except Exception:
logger.warning("LLM evaluation failed, using basic assessment", exc_info=True)
eval_result = {
"score": 0.5,
"assessment": "LLM evaluation unavailable. Manual review required.",
"risks_identified": [],
"recommendation": "needs_review",
"test_scenarios": [],
}
proposal.evaluation_result = eval_result
proposal.evaluated_at = datetime.now(UTC)
proposal.status = "evaluated"
await db.flush()
return {
"proposal_id": str(proposal.id),
"status": proposal.status,
"evaluation": eval_result,
}
# ──────────────────────────────────────────────────────────────────────────
# J-APPROVAL: Human Approval (uses existing approval system)
# ──────────────────────────────────────────────────────────────────────────
async def request_approval(
db: AsyncSession,
tenant_id: uuid.UUID,
proposal_id: uuid.UUID,
requested_by: uuid.UUID,
approver_id: uuid.UUID | None = None,
) -> dict[str, Any]:
"""Create an approval request for a proposal using the existing approval system."""
from app.core.approval import create_approval_request
result = await db.execute(
select(ImprovementProposal).where(
ImprovementProposal.tenant_id == tenant_id,
ImprovementProposal.id == proposal_id,
)
)
proposal = result.scalar_one_or_none()
if not proposal:
return {"error": "Proposal not found"}
if proposal.status not in ("evaluated", "draft"):
return {"error": f"Proposal must be in 'evaluated' or 'draft' status, got '{proposal.status}'"}
req = await create_approval_request(
db=db,
tenant_id=tenant_id,
entity_type="improvement_proposal",
entity_id=proposal.id,
action="activate",
requested_by=requested_by,
requested_by_type="system",
approver_id=approver_id,
metadata={
"proposal_title": proposal.title,
"target_type": proposal.target_type,
"target_name": proposal.target_name,
"evaluation_score": proposal.evaluation_result.get("score", 0.0),
"expected_benefit": proposal.expected_benefit,
},
)
proposal.approval_request_id = req.id
await db.flush()
# Post to Communication via KommunikationContract
try:
from app.plugins.builtins.contracts import get_contract
# Find or create a system conversation for improvement proposals
_komm = get_contract("kommunikation")
if not _komm or not hasattr(_komm, "create_plugin_room"):
from app.plugins.builtins.kommunikation.services import create_plugin_room
room = await create_plugin_room(
db=db, tenant_id=tenant_id, user_id=requested_by,
plugin_name="self_improvement", title="Improvement Proposals",
participant_type="system",
)
conversation_id = room.get("conversation_id") if isinstance(room, dict) else None
if not conversation_id and hasattr(room, "id"):
conversation_id = room.id
if conversation_id:
await KommunikationContract.send_message(
db=db,
tenant_id=tenant_id,
conversation_id=conversation_id,
sender_id=requested_by,
sender_type="system",
content=f"Improvement Proposal: {proposal.title}",
blocks=[
{
"type": "action_card",
"title": f"Improvement Proposal: {proposal.title}",
"content": proposal.description,
"actions": [
{"label": "Approve", "action": "approve", "proposal_id": str(proposal.id), "approval_id": str(req.id)},
{"label": "Reject", "action": "reject", "proposal_id": str(proposal.id), "approval_id": str(req.id)},
],
"metadata": {
"target_type": proposal.target_type,
"target_name": proposal.target_name,
"evaluation_score": proposal.evaluation_result.get("score", 0.0),
"expected_benefit": proposal.expected_benefit,
},
}
],
)
except Exception:
logger.warning("Failed to post proposal to Communication", exc_info=True)
return {
"proposal_id": str(proposal.id),
"approval_request_id": str(req.id),
"status": "pending_approval",
}
# ──────────────────────────────────────────────────────────────────────────
# J-ACTIVATE: Controlled Activation + Rollback
# ──────────────────────────────────────────────────────────────────────────
async def activate_proposal(
db: AsyncSession,
tenant_id: uuid.UUID,
proposal_id: uuid.UUID,
approved_by: uuid.UUID,
) -> dict[str, Any]:
"""Activate an approved proposal — apply the proposed config to the target.
Only applies to agent configurations (system_prompt, tools, temperature, etc.).
Code/plugin patches go through the normal engineering way (J-CODE).
"""
result = await db.execute(
select(ImprovementProposal).where(
ImprovementProposal.tenant_id == tenant_id,
ImprovementProposal.id == proposal_id,
)
)
proposal = result.scalar_one_or_none()
if not proposal:
return {"error": "Proposal not found"}
if proposal.status not in ("approved", "evaluated"):
return {"error": f"Proposal must be 'approved' or 'evaluated', got '{proposal.status}'"}
# Apply config change to target
applied = False
if proposal.target_type == "agent" and proposal.target_ref_id and proposal.proposed_config:
try:
from app.plugins.builtins.automation.models import AgentDefinition
agent_result = await db.execute(
select(AgentDefinition).where(
AgentDefinition.tenant_id == tenant_id,
AgentDefinition.id == proposal.target_ref_id,
)
)
agent = agent_result.scalar_one_or_none()
if agent:
# Create a version snapshot before applying (reuse existing versioning)
from app.plugins.builtins.automation.models import AgentVersion
version_result = await db.execute(
select(func.max(AgentVersion.version_number))
.where(AgentVersion.tenant_id == tenant_id, AgentVersion.agent_id == agent.id)
)
max_ver = version_result.scalar() or 0
version = AgentVersion(
tenant_id=tenant_id,
agent_id=agent.id,
version_number=max_ver + 1,
snapshot={
"system_prompt": agent.system_prompt,
"tool_ids": agent.tool_ids,
"temperature": agent.temperature,
"max_tokens": agent.max_tokens,
"max_steps": agent.max_steps,
},
changed_by=approved_by,
)
db.add(version)
# Apply proposed config
cfg = proposal.proposed_config
if "system_prompt" in cfg:
agent.system_prompt = cfg["system_prompt"]
if "tool_ids" in cfg:
agent.tool_ids = cfg["tool_ids"]
if "temperature" in cfg:
agent.temperature = cfg["temperature"]
if "max_tokens" in cfg:
agent.max_tokens = cfg["max_tokens"]
if "max_steps" in cfg:
agent.max_steps = cfg["max_steps"]
applied = True
except Exception:
logger.exception("Failed to apply proposal to agent")
proposal.status = "active"
proposal.approved_by = approved_by
proposal.approved_at = datetime.now(UTC)
proposal.activated_at = datetime.now(UTC)
await db.flush()
return {
"proposal_id": str(proposal.id),
"status": proposal.status,
"applied": applied,
"activated_at": proposal.activated_at.isoformat() if proposal.activated_at else None,
}
async def rollback_proposal(
db: AsyncSession,
tenant_id: uuid.UUID,
proposal_id: uuid.UUID,
reason: str = "",
) -> dict[str, Any]:
"""Rollback an active proposal — restore previous config."""
result = await db.execute(
select(ImprovementProposal).where(
ImprovementProposal.tenant_id == tenant_id,
ImprovementProposal.id == proposal_id,
)
)
proposal = result.scalar_one_or_none()
if not proposal:
return {"error": "Proposal not found"}
if proposal.status != "active":
return {"error": f"Proposal must be 'active' to rollback, got '{proposal.status}'"}
# Restore previous config
restored = False
if proposal.target_type == "agent" and proposal.target_ref_id and proposal.previous_config:
try:
from app.plugins.builtins.automation.models import AgentDefinition
agent_result = await db.execute(
select(AgentDefinition).where(
AgentDefinition.tenant_id == tenant_id,
AgentDefinition.id == proposal.target_ref_id,
)
)
agent = agent_result.scalar_one_or_none()
if agent:
cfg = proposal.previous_config
if "system_prompt" in cfg:
agent.system_prompt = cfg["system_prompt"]
if "tool_ids" in cfg:
agent.tool_ids = cfg["tool_ids"]
if "temperature" in cfg:
agent.temperature = cfg["temperature"]
if "max_tokens" in cfg:
agent.max_tokens = cfg["max_tokens"]
if "max_steps" in cfg:
agent.max_steps = cfg["max_steps"]
restored = True
except Exception:
logger.exception("Failed to rollback agent config")
proposal.status = "rolled_back"
proposal.rolled_back_at = datetime.now(UTC)
proposal.rollback_reason = reason
await db.flush()
return {
"proposal_id": str(proposal.id),
"status": proposal.status,
"restored": restored,
"rolled_back_at": proposal.rolled_back_at.isoformat() if proposal.rolled_back_at else None,
}
# ──────────────────────────────────────────────────────────────────────────
# J-MEASURE: Impact Measurement
# ──────────────────────────────────────────────────────────────────────────
async def measure_impact(
db: AsyncSession,
tenant_id: uuid.UUID,
proposal_id: uuid.UUID,
) -> dict[str, Any]:
"""Measure pre/post impact of an activated proposal.
Compares agent run metrics before and after activation.
"""
result = await db.execute(
select(ImprovementProposal).where(
ImprovementProposal.tenant_id == tenant_id,
ImprovementProposal.id == proposal_id,
)
)
proposal = result.scalar_one_or_none()
if not proposal:
return {"error": "Proposal not found"}
if proposal.status not in ("active", "rolled_back"):
return {"error": "Proposal must be 'active' or 'rolled_back' to measure"}
# Collect pre-activation metrics (7 days before activation)
pre_metrics: dict[str, Any] = {}
post_metrics: dict[str, Any] = {}
delta: dict[str, Any] = {}
if proposal.target_type == "agent" and proposal.target_ref_id:
try:
from app.plugins.builtins.automation.models import AgentRun
activated = proposal.activated_at or datetime.now(UTC)
pre_start = activated - timedelta(days=7)
# Pre-activation metrics
pre_result = await db.execute(
select(
func.count(AgentRun.id).label("total_runs"),
func.count(AgentRun.id).filter(AgentRun.status.in_(["failed", "error", "timeout"])).label("failed_runs"),
).where(
AgentRun.tenant_id == tenant_id,
AgentRun.agent_id == proposal.target_ref_id,
AgentRun.created_at >= pre_start,
AgentRun.created_at < activated,
)
)
pre_row = pre_result.one_or_none()
pre_total = pre_row.total_runs if pre_row else 0
pre_failed = pre_row.failed_runs if pre_row else 0
pre_metrics = {
"total_runs": pre_total,
"failed_runs": pre_failed,
"failure_rate": (pre_failed / pre_total) if pre_total > 0 else 0.0,
}
# Post-activation metrics
post_result = await db.execute(
select(
func.count(AgentRun.id).label("total_runs"),
func.count(AgentRun.id).filter(AgentRun.status.in_(["failed", "error", "timeout"])).label("failed_runs"),
).where(
AgentRun.tenant_id == tenant_id,
AgentRun.agent_id == proposal.target_ref_id,
AgentRun.created_at >= activated,
)
)
post_row = post_result.one_or_none()
post_total = post_row.total_runs if post_row else 0
post_failed = post_row.failed_runs if post_row else 0
post_metrics = {
"total_runs": post_total,
"failed_runs": post_failed,
"failure_rate": (post_failed / post_total) if post_total > 0 else 0.0,
}
# Compute delta
delta = {
"total_runs_change": post_total - pre_total,
"failure_rate_change": post_metrics["failure_rate"] - pre_metrics["failure_rate"],
}
except Exception:
logger.warning("Failed to collect agent run metrics", exc_info=True)
# Determine if impact is positive
failure_rate_improved = delta.get("failure_rate_change", 0) < 0
is_positive = "positive" if failure_rate_improved else ("neutral" if delta.get("failure_rate_change", 0) == 0 else "negative")
assessment = f"Failure rate changed from {pre_metrics.get('failure_rate', 0):.1%} to {post_metrics.get('failure_rate', 0):.1%}. "
assessment += "Improvement detected." if failure_rate_improved else "No improvement or regression detected."
measurement = ImpactMeasurement(
tenant_id=tenant_id,
proposal_id=proposal_id,
pre_metrics=pre_metrics,
post_metrics=post_metrics,
delta=delta,
assessment=assessment,
is_positive=is_positive,
)
db.add(measurement)
await db.flush()
return {
"measurement_id": str(measurement.id),
"proposal_id": str(proposal_id),
"pre_metrics": pre_metrics,
"post_metrics": post_metrics,
"delta": delta,
"assessment": assessment,
"is_positive": is_positive,
}
# ──────────────────────────────────────────────────────────────────────────
# Query helpers
# ──────────────────────────────────────────────────────────────────────────
async def list_signals(
db: AsyncSession,
tenant_id: uuid.UUID,
page: int = 1,
page_size: int = 20,
source_type: str | None = None,
) -> dict[str, Any]:
"""List improvement signals with pagination."""
q = select(ImprovementSignal).where(
ImprovementSignal.tenant_id == tenant_id,
)
if source_type:
q = q.where(ImprovementSignal.source_type == source_type)
q = q.order_by(ImprovementSignal.created_at.desc())
count_q = select(func.count()).select_from(q.subquery())
total = (await db.execute(count_q)).scalar() or 0
offset = (page - 1) * page_size
result = await db.execute(q.offset(offset).limit(page_size))
items = [
{
"id": str(s.id),
"source_type": s.source_type,
"signal_kind": s.signal_kind,
"severity": s.severity,
"summary": s.summary,
"confidence": s.confidence,
"pattern_id": str(s.pattern_id) if s.pattern_id else None,
"created_at": s.created_at.isoformat() if s.created_at else None,
}
for s in result.scalars().all()
]
return {"items": items, "total": total, "page": page, "page_size": page_size}
async def list_patterns(
db: AsyncSession,
tenant_id: uuid.UUID,
page: int = 1,
page_size: int = 20,
) -> dict[str, Any]:
"""List detected patterns with pagination."""
q = select(ImprovementPattern).where(
ImprovementPattern.tenant_id == tenant_id,
).order_by(ImprovementPattern.created_at.desc())
count_q = select(func.count()).select_from(q.subquery())
total = (await db.execute(count_q)).scalar() or 0
offset = (page - 1) * page_size
result = await db.execute(q.offset(offset).limit(page_size))
items = [
{
"id": str(p.id),
"pattern_kind": p.pattern_kind,
"title": p.title,
"description": p.description,
"target_type": p.target_type,
"target_name": p.target_name,
"occurrence_count": p.occurrence_count,
"confidence": p.confidence,
"status": p.status,
"proposed_action": p.proposed_action,
"created_at": p.created_at.isoformat() if p.created_at else None,
}
for p in result.scalars().all()
]
return {"items": items, "total": total, "page": page, "page_size": page_size}
async def list_proposals(
db: AsyncSession,
tenant_id: uuid.UUID,
page: int = 1,
page_size: int = 20,
status: str | None = None,
) -> dict[str, Any]:
"""List improvement proposals with pagination."""
q = select(ImprovementProposal).where(
ImprovementProposal.tenant_id == tenant_id,
)
if status:
q = q.where(ImprovementProposal.status == status)
q = q.order_by(ImprovementProposal.created_at.desc())
count_q = select(func.count()).select_from(q.subquery())
total = (await db.execute(count_q)).scalar() or 0
offset = (page - 1) * page_size
result = await db.execute(q.offset(offset).limit(page_size))
items = [
{
"id": str(p.id),
"title": p.title,
"description": p.description,
"target_type": p.target_type,
"target_name": p.target_name,
"status": p.status,
"version_number": p.version_number,
"rationale": p.rationale,
"expected_benefit": p.expected_benefit,
"risk_assessment": p.risk_assessment,
"evaluation_score": p.evaluation_result.get("score", 0.0) if p.evaluation_result else 0.0,
"pattern_id": str(p.pattern_id) if p.pattern_id else None,
"created_at": p.created_at.isoformat() if p.created_at else None,
"activated_at": p.activated_at.isoformat() if p.activated_at else None,
}
for p in result.scalars().all()
]
return {"items": items, "total": total, "page": page, "page_size": page_size}
async def get_proposal_detail(
db: AsyncSession,
tenant_id: uuid.UUID,
proposal_id: uuid.UUID,
) -> dict[str, Any]:
"""Get full proposal detail."""
result = await db.execute(
select(ImprovementProposal).where(
ImprovementProposal.tenant_id == tenant_id,
ImprovementProposal.id == proposal_id,
)
)
p = result.scalar_one_or_none()
if not p:
return {"error": "Proposal not found"}
return {
"id": str(p.id),
"title": p.title,
"description": p.description,
"target_type": p.target_type,
"target_ref_id": str(p.target_ref_id) if p.target_ref_id else None,
"target_name": p.target_name,
"proposed_config": p.proposed_config,
"previous_config": p.previous_config,
"evidence_refs": p.evidence_refs,
"rationale": p.rationale,
"expected_benefit": p.expected_benefit,
"risk_assessment": p.risk_assessment,
"status": p.status,
"version_number": p.version_number,
"evaluation_result": p.evaluation_result,
"evaluated_at": p.evaluated_at.isoformat() if p.evaluated_at else None,
"approval_request_id": str(p.approval_request_id) if p.approval_request_id else None,
"approved_by": str(p.approved_by) if p.approved_by else None,
"approved_at": p.approved_at.isoformat() if p.approved_at else None,
"activated_at": p.activated_at.isoformat() if p.activated_at else None,
"rolled_back_at": p.rolled_back_at.isoformat() if p.rolled_back_at else None,
"rollback_reason": p.rollback_reason,
"pattern_id": str(p.pattern_id) if p.pattern_id else None,
"created_at": p.created_at.isoformat() if p.created_at else None,
}