File size: 2,526 Bytes
cb505ff
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
from __future__ import annotations

from typing import Any, Dict, Optional
from src.skills.models import SkillDefinition, SkillResult
from src.skills.registry import SkillRegistry
from src.runtime.agent_runtime import AgentRuntime
from src.runtime.execution_context import ExecutionContext
from src.skills.adapters import LegacyAgentSkillAdapter
import time


class WorkerRuntime:
    """Ephemeral Worker Runtime slot for short-lived task execution."""

    def __init__(self, registry: SkillRegistry):
        self.registry = registry
        self.agent_runtime = AgentRuntime(name="WorkerRuntime", registry=registry)
        self.is_destroyed = False
        self.created_at = time.perf_counter()

    @property
    def active_skill(self) -> Optional[SkillDefinition]:
        return self.agent_runtime.active_skill

    def load_skill(self, skill_id: str, version: Optional[str] = None) -> SkillDefinition:
        if self.is_destroyed:
            raise RuntimeError("WorkerRuntime has been destroyed. Create a new instance.")
        config = self.agent_runtime.load_skill(skill_id, version=version)
        return config.definition

    def execute(self, context: ExecutionContext, input_override: Optional[Dict[str, Any]] = None) -> SkillResult:
        if self.is_destroyed:
            raise RuntimeError("WorkerRuntime has been destroyed. Create a new instance.")
        result = self.agent_runtime.execute(context=context, input_override=input_override)
        context.worker_result = result.output
        return result

    async def execute_adapter(
        self,
        adapter: LegacyAgentSkillAdapter,
        state: Dict[str, Any],
        config: Optional[Any] = None,
    ) -> SkillResult:
        """Execute one compatibility adapter inside this worker lifecycle."""
        if self.is_destroyed:
            raise RuntimeError("WorkerRuntime has been destroyed. Create a new instance.")
        if self.active_skill is None or self.active_skill.id != adapter.skill_id:
            self.load_skill(adapter.skill_id, version=adapter.definition.version)
        return await adapter.execute(state, config=config)

    def flush_and_destroy(self) -> None:
        """Flush active skill state and mark worker as destroyed."""
        self.agent_runtime.flush()
        self.agent_runtime.metrics.record(
            "runtime.destroyed",
            runtime="WorkerRuntime",
            worker_lifetime_ms=round((time.perf_counter() - self.created_at) * 1000, 3),
        )
        self.is_destroyed = True