Download src/runtime/worker_runtime.py from DiabetesCareChatbot/dmChatbotBackend: direct link, hf CLI and curl.
- Browser
- Download file 2.53 kB
-
https://huggingface.co/spaces/DiabetesCareChatbot/dmChatbotBackend/resolve/main/src/runtime/worker_runtime.py
- Command line
-
hf download hf://spaces/DiabetesCareChatbot/dmChatbotBackend/src/runtime/worker_runtime.py
-
curl -L -o worker_runtime.py https://huggingface.co/spaces/DiabetesCareChatbot/dmChatbotBackend/resolve/main/src/runtime/worker_runtime.py
2.53 kB
| 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() | |
| 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 | |