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
|