1059 lines
41 KiB
Python
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 (
|
|
ImpactMeasurement,
|
|
ImprovementPattern,
|
|
ImprovementProposal,
|
|
ImprovementSignal,
|
|
)
|
|
|
|
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 (ARCH-013: contract
|
|
# only — no direct plugin imports; skip cleanly when contract is absent)
|
|
try:
|
|
from app.plugins.builtins.contracts import get_contract
|
|
_komm = get_contract("kommunikation")
|
|
if not _komm or not hasattr(_komm, "create_plugin_room"):
|
|
logger.warning("kommunikation contract unavailable - skipping proposal notification")
|
|
else:
|
|
room = await _komm.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 _komm.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,
|
|
}
|