Implement intelligent agent learning from Knowledge Graph execution history with per-task-type expertise tracking, recency bias, and learning curves. ## Phase 5.3 Implementation ### Learning Infrastructure (✅ Complete) - LearningProfileService with per-task-type expertise metrics - TaskTypeExpertise model tracking success_rate, confidence, learning curves - Recency bias weighting: recent 7 days weighted 3x higher (exponential decay) - Confidence scoring prevents overfitting: min(1.0, executions / 20) - Learning curves computed from daily execution windows ### Agent Scoring Service (✅ Complete) - Unified AgentScore combining SwarmCoordinator + learning profiles - Scoring formula: 0.3*base + 0.5*expertise + 0.2*confidence - Rank agents by combined score for intelligent assignment - Support for recency-biased scoring (recent_success_rate) - Methods: rank_agents, select_best, rank_agents_with_recency ### KG Integration (✅ Complete) - KGPersistence::get_executions_for_task_type() - query by agent + task type - KGPersistence::get_agent_executions() - all executions for agent - Coordinator::load_learning_profile_from_kg() - core KG→Learning integration - Coordinator::load_all_learning_profiles() - batch load for multiple agents - Convert PersistedExecution → ExecutionData for learning calculations ### Agent Assignment Integration (✅ Complete) - AgentCoordinator uses learning profiles for task assignment - extract_task_type() infers task type from title/description - assign_task() scores candidates using AgentScoringService - Fallback to load-based selection if no learning data available - Learning profiles stored in coordinator.learning_profiles RwLock ### Profile Adapter Enhancements (✅ Complete) - create_learning_profile() - initialize empty profiles - add_task_type_expertise() - set task-type expertise - update_profile_with_learning() - update swarm profiles from learning ## Files Modified ### vapora-knowledge-graph/src/persistence.rs (+30 lines) - get_executions_for_task_type(agent_id, task_type, limit) - get_agent_executions(agent_id, limit) ### vapora-agents/src/coordinator.rs (+100 lines) - load_learning_profile_from_kg() - core KG integration method - load_all_learning_profiles() - batch loading for agents - assign_task() already uses learning-based scoring via AgentScoringService ### Existing Complete Implementation - vapora-knowledge-graph/src/learning.rs - calculation functions - vapora-agents/src/learning_profile.rs - data structures and expertise - vapora-agents/src/scoring.rs - unified scoring service - vapora-agents/src/profile_adapter.rs - adapter methods ## Tests Passing - learning_profile: 7 tests ✅ - scoring: 5 tests ✅ - profile_adapter: 6 tests ✅ - coordinator: learning-specific tests ✅ ## Data Flow 1. Task arrives → AgentCoordinator::assign_task() 2. Extract task_type from description 3. Query KG for task-type executions (load_learning_profile_from_kg) 4. Calculate expertise with recency bias 5. Score candidates (SwarmCoordinator + learning) 6. Assign to top-scored agent 7. Execution result → KG → Update learning profiles ## Key Design Decisions ✅ Recency bias: 7-day half-life with 3x weight for recent performance ✅ Confidence scoring: min(1.0, total_executions / 20) prevents overfitting ✅ Hierarchical scoring: 30% base load, 50% expertise, 20% confidence ✅ KG query limit: 100 recent executions per task-type for performance ✅ Async loading: load_learning_profile_from_kg supports concurrent loads ## Next: Phase 5.4 - Cost Optimization Ready to implement budget enforcement and cost-aware provider selection.
63 lines
3.8 KiB
Plaintext
63 lines
3.8 KiB
Plaintext
-- Migration 003: Workflow Engine
|
|
-- Creates tables for workflow definitions and execution tracking
|
|
|
|
-- Workflows table (workflow definitions)
|
|
DEFINE TABLE workflows SCHEMAFULL
|
|
PERMISSIONS
|
|
FOR select WHERE tenant_id = $auth.tenant_id
|
|
FOR create, update, delete WHERE tenant_id = $auth.tenant_id AND ("admin" IN $auth.roles OR "project_manager" IN $auth.roles);
|
|
|
|
DEFINE FIELD id ON TABLE workflows TYPE record<workflows>;
|
|
DEFINE FIELD tenant_id ON TABLE workflows TYPE string ASSERT $value != NONE;
|
|
DEFINE FIELD name ON TABLE workflows TYPE string ASSERT $value != NONE;
|
|
DEFINE FIELD description ON TABLE workflows TYPE option<string>;
|
|
DEFINE FIELD status ON TABLE workflows TYPE string ASSERT $value INSIDE ["draft", "active", "paused", "completed", "failed"] DEFAULT "draft";
|
|
DEFINE FIELD definition ON TABLE workflows TYPE object DEFAULT {};
|
|
DEFINE FIELD created_at ON TABLE workflows TYPE datetime DEFAULT time::now();
|
|
DEFINE FIELD updated_at ON TABLE workflows TYPE datetime DEFAULT time::now() VALUE time::now();
|
|
|
|
DEFINE INDEX idx_workflows_tenant ON TABLE workflows COLUMNS tenant_id;
|
|
DEFINE INDEX idx_workflows_status ON TABLE workflows COLUMNS status;
|
|
DEFINE INDEX idx_workflows_tenant_status ON TABLE workflows COLUMNS tenant_id, status;
|
|
|
|
-- Workflow steps table (execution tracking)
|
|
DEFINE TABLE workflow_steps SCHEMAFULL
|
|
PERMISSIONS
|
|
FOR select WHERE $parent.tenant_id = $auth.tenant_id
|
|
FOR create, update WHERE $parent.tenant_id = $auth.tenant_id;
|
|
|
|
DEFINE FIELD id ON TABLE workflow_steps TYPE record<workflow_steps>;
|
|
DEFINE FIELD workflow_id ON TABLE workflow_steps TYPE string ASSERT $value != NONE;
|
|
DEFINE FIELD step_id ON TABLE workflow_steps TYPE string ASSERT $value != NONE;
|
|
DEFINE FIELD step_name ON TABLE workflow_steps TYPE string ASSERT $value != NONE;
|
|
DEFINE FIELD agent_id ON TABLE workflow_steps TYPE option<string>;
|
|
DEFINE FIELD status ON TABLE workflow_steps TYPE string ASSERT $value INSIDE ["pending", "in_progress", "completed", "failed", "skipped"] DEFAULT "pending";
|
|
DEFINE FIELD result ON TABLE workflow_steps TYPE option<object>;
|
|
DEFINE FIELD error_message ON TABLE workflow_steps TYPE option<string>;
|
|
DEFINE FIELD started_at ON TABLE workflow_steps TYPE option<datetime>;
|
|
DEFINE FIELD completed_at ON TABLE workflow_steps TYPE option<datetime>;
|
|
DEFINE FIELD created_at ON TABLE workflow_steps TYPE datetime DEFAULT time::now();
|
|
|
|
DEFINE INDEX idx_workflow_steps_workflow ON TABLE workflow_steps COLUMNS workflow_id;
|
|
DEFINE INDEX idx_workflow_steps_status ON TABLE workflow_steps COLUMNS status;
|
|
DEFINE INDEX idx_workflow_steps_agent ON TABLE workflow_steps COLUMNS agent_id;
|
|
DEFINE INDEX idx_workflow_steps_workflow_step ON TABLE workflow_steps COLUMNS workflow_id, step_id UNIQUE;
|
|
|
|
-- Workflow executions table (execution history)
|
|
DEFINE TABLE workflow_executions SCHEMAFULL
|
|
PERMISSIONS
|
|
FOR select WHERE tenant_id = $auth.tenant_id
|
|
FOR create WHERE tenant_id = $auth.tenant_id;
|
|
|
|
DEFINE FIELD id ON TABLE workflow_executions TYPE record<workflow_executions>;
|
|
DEFINE FIELD tenant_id ON TABLE workflow_executions TYPE string ASSERT $value != NONE;
|
|
DEFINE FIELD workflow_id ON TABLE workflow_executions TYPE string ASSERT $value != NONE;
|
|
DEFINE FIELD status ON TABLE workflow_executions TYPE string ASSERT $value INSIDE ["running", "completed", "failed", "cancelled"] DEFAULT "running";
|
|
DEFINE FIELD started_at ON TABLE workflow_executions TYPE datetime DEFAULT time::now();
|
|
DEFINE FIELD completed_at ON TABLE workflow_executions TYPE option<datetime>;
|
|
DEFINE FIELD duration_ms ON TABLE workflow_executions TYPE option<int>;
|
|
|
|
DEFINE INDEX idx_workflow_executions_workflow ON TABLE workflow_executions COLUMNS workflow_id;
|
|
DEFINE INDEX idx_workflow_executions_tenant ON TABLE workflow_executions COLUMNS tenant_id;
|
|
DEFINE INDEX idx_workflow_executions_status ON TABLE workflow_executions COLUMNS status;
|