File size: 6,617 Bytes
ff0e46c
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
"""

Luuna Autoflow Core - Execution Bus

===================================



Central execution orchestrator that:

- Creates execution IDs

- Calls router for adapter selection

- Calls adapter for execution

- Catches exceptions

- Always returns JSON

- Logs status

"""

import logging
from typing import Dict, Optional
from datetime import datetime
import uuid

from .models import (
    AutoflowTask,
    AutoflowResult,
    ExecutionRecord,
    TaskStatus,
)
from .router import Router
from .policy import PolicyEngine
from .memory import MemoryStore
from .adapters.base import BaseAdapter

logger = logging.getLogger(__name__)


class ExecutionBus:
    """

    Central execution orchestrator for Luuna Autoflow.

    

    All executions go through this bus for:

    - Auditing

    - Policy enforcement

    - Error handling

    - Result formatting

    """
    
    def __init__(

        self, 

        adapters: Dict[str, BaseAdapter],

        memory: Optional[MemoryStore] = None,

        policy: Optional[PolicyEngine] = None,

    ):
        self.adapters = adapters
        self.memory = memory or MemoryStore()
        self.policy = policy or PolicyEngine()
        self.router = Router()
    
    def execute(self, task: AutoflowTask) -> AutoflowResult:
        """

        Execute a task through the appropriate adapter.

        

        Args:

            task: The task to execute

            

        Returns:

            AutoflowResult with execution outcome

        """
        # Create execution ID
        execution_id = str(uuid.uuid4())
        
        # Create initial record
        record = ExecutionRecord(
            execution_id=execution_id,
            goal=task.goal,
            domain=task.domain,
            mode=task.mode,
            status=TaskStatus.PENDING,
            requires_approval=task.approval_required,
        )
        self.memory.store(record)
        
        logger.info(f"[Autoflow] Starting execution {execution_id}: {task.goal[:100]}...")
        
        try:
            # Policy check
            policy_result = self.policy.check(task)
            if not policy_result.allowed:
                return self._create_blocked_result(
                    execution_id, 
                    record, 
                    policy_result.reason
                )
            
            # Route to adapter
            try:
                adapter_id, adapter = self.router.select_adapter(task, self.adapters)
            except ValueError as e:
                return self._create_error_result(
                    execution_id, 
                    record, 
                    str(e)
                )
            
            record.selected_adapter = adapter_id
            record.status = TaskStatus.RUNNING
            self.memory.update(record)
            
            logger.info(f"[Autoflow] Routed to adapter: {adapter_id}")
            
            # Generate plan
            plan = adapter.plan(task)
            record.plan = plan
            
            # Check if approval required
            requires_approval = (
                task.approval_required
                or policy_result.requires_approval
                or adapter.capabilities.requires_approval
            )
            
            # Execute based on mode
            result_data = {}
            warnings = []
            
            if task.mode.value == "execute_mock":
                # Only execute in mock mode
                if adapter.can_handle(task):
                    result_data = adapter.execute(task)
                    warnings.append("Executed in mock mode - no real actions taken")
                else:
                    warnings.append("Adapter cannot handle task - plan only")
            else:
                warnings.append("Plan-only mode - no execution performed")
            
            # Update record
            record.status = TaskStatus.COMPLETED
            record.result = result_data
            record.warnings = warnings
            record.requires_approval = requires_approval
            record.completed_at = datetime.utcnow()
            self.memory.update(record)
            
            logger.info(f"[Autoflow] Execution {execution_id} completed successfully")
            
            return AutoflowResult(
                success=True,
                execution_id=execution_id,
                selected_adapter=adapter_id,
                plan=plan,
                result=result_data,
                warnings=warnings,
                requires_approval=requires_approval,
                status=TaskStatus.COMPLETED,
            )
            
        except Exception as e:
            logger.error(f"[Autoflow] Execution {execution_id} failed: {str(e)}")
            return self._create_error_result(execution_id, record, str(e))
    
    def get_execution(self, execution_id: str) -> Optional[ExecutionRecord]:
        """Retrieve an execution record by ID."""
        return self.memory.get(execution_id)
    
    def _create_error_result(

        self, 

        execution_id: str, 

        record: ExecutionRecord, 

        error: str

    ) -> AutoflowResult:
        """Create an error result."""
        record.status = TaskStatus.FAILED
        record.warnings = [error]
        record.completed_at = datetime.utcnow()
        self.memory.update(record)
        
        return AutoflowResult(
            success=False,
            execution_id=execution_id,
            selected_adapter=record.selected_adapter or "none",
            plan=record.plan,
            result={"error": error},
            warnings=[error],
            requires_approval=False,
            status=TaskStatus.FAILED,
        )
    
    def _create_blocked_result(

        self, 

        execution_id: str, 

        record: ExecutionRecord, 

        reason: str

    ) -> AutoflowResult:
        """Create a blocked result from policy."""
        record.status = TaskStatus.REQUIRES_APPROVAL
        record.warnings = [reason]
        record.requires_approval = True
        self.memory.update(record)
        
        return AutoflowResult(
            success=False,
            execution_id=execution_id,
            selected_adapter="none",
            plan=[],
            result={"blocked": True, "reason": reason},
            warnings=[reason],
            requires_approval=True,
            status=TaskStatus.REQUIRES_APPROVAL,
        )