Spaces:
Running on Zero
Running on Zero
| from __future__ import annotations | |
| from collections.abc import Callable | |
| from contextlib import ExitStack | |
| from dataclasses import dataclass | |
| from threading import Lock | |
| from pathlib import Path | |
| from typing import Final, TypeAlias, final, override | |
| import anyio | |
| import os | |
| from anyio import to_thread | |
| from anyio.from_thread import BlockingPortal, start_blocking_portal | |
| from pydantic import BaseModel, JsonValue | |
| from .event import ( | |
| BacklogItemData, | |
| EventModel, | |
| EventEnvelope, | |
| GenerateData, | |
| TextChunk, | |
| TodoItemData, | |
| VoiceStatus, | |
| VoiceStatusData, | |
| VoiceDelegationData, | |
| WorkflowCompletedData, | |
| WorkflowCheckpointResult, | |
| WorkflowPhaseCompleteData, | |
| WorkflowPhaseStartData, | |
| WorkflowStartedData, | |
| ) | |
| from .manager import PluginManager | |
| from .plugin import Plugin | |
| from .contracts import ProviderSelection | |
| from .providers import ( | |
| FixtureProvider, | |
| SpaceProvider, | |
| SpaceProviderConfig, | |
| build_provider, | |
| ) | |
| JsonRecord: TypeAlias = dict[str, JsonValue] | |
| RuntimeBuilder: TypeAlias = Callable[ | |
| [ProviderSelection], tuple[PluginManager, tuple[Plugin, ...]] | |
| ] | |
| _FIXTURE_WORKFLOW: Final[str] = "space-plan" | |
| class SpacePlanData(EventModel): | |
| prompt: str | |
| provider: str | |
| voice: str | |
| plan: str | |
| model_response: str | |
| class SpaceArtifactData(EventModel): | |
| name: str | |
| kind: str | |
| content: str | |
| path: str | |
| class RuntimeSnapshot: | |
| provider: str | |
| model: str | |
| mode: str | |
| closed: bool | |
| events: tuple[JsonRecord, ...] | |
| notifications: tuple[JsonRecord, ...] | |
| plans: tuple[JsonRecord, ...] | |
| agent_statuses: tuple[JsonRecord, ...] | |
| workflows: tuple[JsonRecord, ...] | |
| mcp: tuple[JsonRecord, ...] | |
| plugins: tuple[JsonRecord, ...] | |
| artifacts: tuple[JsonRecord, ...] | |
| def as_outputs(self) -> JsonRecord: | |
| from .presentation import snapshot_outputs | |
| return snapshot_outputs(self) | |
| class _FixtureRuntimePlugin(Plugin): | |
| _selection: ProviderSelection | |
| _provider: SpaceProvider | |
| def __init__(self, selection: ProviderSelection, provider: SpaceProvider) -> None: | |
| super().__init__("space_fixture_runtime") | |
| self._selection = selection | |
| self._provider = provider | |
| async def initialize(self) -> None: | |
| manager = self.manager | |
| if manager is None: | |
| raise RuntimeError("Space fixture runtime requires a manager") | |
| manager.register_event("space_plan", SpacePlanData) | |
| manager.register_event("space_artifact", SpaceArtifactData) | |
| manager.register_event("voice_status", VoiceStatusData) | |
| manager.register_event("voice_delegation", VoiceDelegationData) | |
| manager.register_event("backlog_item", BacklogItemData) | |
| manager.register_event("workflow_started", WorkflowStartedData) | |
| manager.register_event("workflow_phase_start", WorkflowPhaseStartData) | |
| manager.register_event("workflow_phase_complete", WorkflowPhaseCompleteData) | |
| manager.register_event("workflow_completed", WorkflowCompletedData) | |
| manager.register_event("todo_updated", TodoItemData) | |
| async def on_generate(self, data: GenerateData) -> None: | |
| prompt = data.prompt | |
| if prompt is None: | |
| return | |
| completion = await self._provider.complete( | |
| "Create a concise project plan with an Objective section and actionable " | |
| "checkbox steps for this request:\n\n" + prompt | |
| ) | |
| manager = self.manager | |
| if manager is None: | |
| raise RuntimeError("Space fixture runtime is not attached") | |
| _ = await manager.emit( | |
| "voice_status", | |
| VoiceStatusData(voice="jane", status=VoiceStatus.THINKING), | |
| source=self.name, | |
| ) | |
| _ = await manager.emit( | |
| "voice_status", | |
| VoiceStatusData(voice="doktor", status=VoiceStatus.THINKING), | |
| source=self.name, | |
| ) | |
| _ = await manager.emit( | |
| "voice_status", | |
| VoiceStatusData(voice="conquest", status=VoiceStatus.WAITING), | |
| source=self.name, | |
| ) | |
| items = _plan_items(completion, prompt) | |
| _ = await manager.emit( | |
| "voice_delegation", | |
| VoiceDelegationData(from_voice="jane", to_voice="doktor", task="Draft PLAN.md", result=completion), | |
| source=self.name, | |
| ) | |
| _ = await manager.emit( | |
| "space_plan", | |
| SpacePlanData( | |
| prompt=prompt, | |
| provider=self._selection.provider, | |
| voice="doktor", | |
| plan=f"Jane delegated planning to Doktor for: {prompt}", | |
| model_response=completion, | |
| ), | |
| source=self.name, | |
| ) | |
| storage_dir = Path(os.getenv("KATETO_SPACE_STORAGE_DIR", "/tmp/kateto-space")) | |
| storage_dir.mkdir(parents=True, exist_ok=True) | |
| plan_content = f"# PLAN.md\n\n## Objective\n\n{prompt}\n\n## Model plan\n\n{completion.strip()}\n" | |
| backlog_content = "# BACKLOG.md\n\n" + "\n".join( | |
| f"- [ ] **{item['priority']}** {item['item']} _(owner: {item['owner']})_" | |
| for item in items | |
| ) + "\n" | |
| todo_content = "# TODO.md\n\n" + "\n".join( | |
| f"- [ ] {item['item']} _(owner: {item['owner']})_" for item in items | |
| ) + "\n" | |
| workflow_content = ( | |
| "# space-plan\n\n" | |
| "Plan → PLAN.md → Triangulate → BACKLOG.md → Update → TODO.md.\n" | |
| ) | |
| artifacts = ( | |
| ("PLAN.md", "plan", plan_content), | |
| ("TODO.md", "todo", todo_content), | |
| ("BACKLOG.md", "backlog", backlog_content), | |
| ("workflows/space-plan/workflow.md", "workflow", workflow_content), | |
| ("voices/jane/SOUL.md", "soul", "# Jane\n\nOrchestrator: delegates planning, triangulation, and updates.\n"), | |
| ("voices/doktor/SOUL.md", "soul", "# Doktor\n\nPlanner: turns requests into executable plans and backlog items.\n"), | |
| ("voices/conquest/SOUL.md", "soul", "# Conquest\n\nTriangulator: challenges assumptions and checks delivery readiness.\n"), | |
| ) | |
| for name, kind, content in artifacts: | |
| path = storage_dir / name | |
| path.parent.mkdir(parents=True, exist_ok=True) | |
| path.write_text(content, encoding="utf-8") | |
| _ = await manager.emit( | |
| "space_artifact", | |
| SpaceArtifactData(name=name, kind=kind, content=content, path=str(path)), | |
| source=self.name, | |
| ) | |
| _ = await manager.emit( | |
| "workflow_started", | |
| WorkflowStartedData( | |
| workflow=_FIXTURE_WORKFLOW, voice="jane", context={"prompt": prompt} | |
| ), | |
| source=self.name, | |
| ) | |
| _ = await manager.emit( | |
| "workflow_phase_start", | |
| WorkflowPhaseStartData( | |
| workflow=_FIXTURE_WORKFLOW, | |
| phase_id="plan", | |
| voice="doktor", | |
| instructions=["Turn the request into an actionable plan"], | |
| ), | |
| source=self.name, | |
| ) | |
| _ = await manager.emit( | |
| "workflow_phase_complete", | |
| WorkflowPhaseCompleteData( | |
| workflow=_FIXTURE_WORKFLOW, | |
| phase_id="plan", | |
| voice="doktor", | |
| deliverables=["PLAN.md"], | |
| checkpoint_results=[ | |
| WorkflowCheckpointResult(checkpoint="plan-recorded", passed=True), | |
| ], | |
| ), | |
| source=self.name, | |
| ) | |
| _ = await manager.emit( | |
| "voice_delegation", | |
| VoiceDelegationData(from_voice="jane", to_voice="conquest", task="Triangulate PLAN.md into backlog items", result=f"Triangulated {len(items)} actionable items."), | |
| source=self.name, | |
| ) | |
| _ = await manager.emit("voice_status", VoiceStatusData(voice="conquest", status=VoiceStatus.THINKING), source=self.name) | |
| _ = await manager.emit("workflow_phase_start", WorkflowPhaseStartData(workflow=_FIXTURE_WORKFLOW, phase_id="triangulate", voice="conquest", instructions=["Convert plan steps into backlog items"]), source=self.name) | |
| for item in items: | |
| _ = await manager.emit("backlog_item", BacklogItemData(voice="conquest", item=item["item"], priority=item["priority"], source="PLAN.md"), source=self.name) | |
| _ = await manager.emit("workflow_phase_complete", WorkflowPhaseCompleteData(workflow=_FIXTURE_WORKFLOW, phase_id="triangulate", voice="conquest", deliverables=["BACKLOG.md"], checkpoint_results=[WorkflowCheckpointResult(checkpoint="backlog-derived", passed=True)]), source=self.name) | |
| _ = await manager.emit("voice_delegation", VoiceDelegationData(from_voice="jane", to_voice="doktor", task="Update TODO.md and voice SOUL files", result="Project records updated."), source=self.name) | |
| _ = await manager.emit("workflow_phase_start", WorkflowPhaseStartData(workflow=_FIXTURE_WORKFLOW, phase_id="update", voice="doktor", instructions=["Persist TODO.md and personality changes"]), source=self.name) | |
| for item in items: | |
| _ = await manager.emit("todo_updated", TodoItemData(voice="doktor", task=item["item"], completed=False), source=self.name) | |
| _ = await manager.emit("workflow_phase_complete", WorkflowPhaseCompleteData(workflow=_FIXTURE_WORKFLOW, phase_id="update", voice="doktor", deliverables=["TODO.md", "voices/*/SOUL.md"], checkpoint_results=[WorkflowCheckpointResult(checkpoint="records-updated", passed=True)]), source=self.name) | |
| _ = await manager.emit( | |
| "workflow_completed", | |
| WorkflowCompletedData(workflow=_FIXTURE_WORKFLOW, voice="jane"), | |
| source=self.name, | |
| ) | |
| _ = await manager.emit( | |
| "text_chunk", | |
| TextChunk(text="Jane completed the orchestration and updated the project files.", sequence=0, final=True, voice_id="jane"), | |
| source=self.name, | |
| ) | |
| _ = await manager.emit( | |
| "voice_status", | |
| VoiceStatusData(voice="doktor", status=VoiceStatus.IDLE), | |
| source=self.name, | |
| ) | |
| for voice in ("jane", "conquest"): | |
| _ = await manager.emit( | |
| "voice_status", | |
| VoiceStatusData(voice=voice, status=VoiceStatus.IDLE), | |
| source=self.name, | |
| ) | |
| def _plan_items(completion: str, prompt: str) -> list[dict[str, str]]: | |
| items: list[dict[str, str]] = [] | |
| source = "" if completion.startswith("Plan ready for:") else completion | |
| for line in source.splitlines(): | |
| text = line.strip().lstrip("-* ").strip() | |
| if text.startswith("[ ]"): | |
| text = text[3:].strip() | |
| if text and len(text) > 8 and not text[0].isdigit(): | |
| items.append({"item": text, "priority": "P1", "owner": "doktor"}) | |
| if not items: | |
| items = [ | |
| {"item": f"Clarify scope and success criteria for: {prompt}", "priority": "P1", "owner": "doktor"}, | |
| {"item": "Execute the planned work and capture decisions", "priority": "P1", "owner": "conquest"}, | |
| {"item": "Validate outcomes and summarize next steps", "priority": "P2", "owner": "jane"}, | |
| ] | |
| return items[:8] | |
| def _fixture_builder( | |
| selection: ProviderSelection, | |
| ) -> tuple[PluginManager, tuple[Plugin, ...]]: | |
| manager = PluginManager(event_limit=200) | |
| return manager, (_FixtureRuntimePlugin(selection, FixtureProvider()),) | |
| def _live_builder( | |
| selection: ProviderSelection, | |
| config: SpaceProviderConfig, | |
| ) -> tuple[PluginManager, tuple[Plugin, ...]]: | |
| manager = PluginManager(event_limit=200) | |
| return manager, ( | |
| _FixtureRuntimePlugin(selection, build_provider(selection, config)), | |
| ) | |
| class SpaceRuntimeSession: | |
| def __init__( | |
| self, | |
| selection: ProviderSelection, | |
| manager: PluginManager, | |
| plugins: tuple[Plugin, ...], | |
| *, | |
| mode: str, | |
| model: str, | |
| ) -> None: | |
| self.provider: str = selection.provider | |
| self.model: str = model | |
| self.mode: str = mode | |
| self.manager: PluginManager = manager | |
| self._plugins: tuple[Plugin, ...] = plugins | |
| self._session_key: str | None = selection.session_key | |
| self._started: bool = False | |
| self._closed: bool = False | |
| self._prompt_lock: anyio.Lock | None = None | |
| self._portal_stack: ExitStack | None = None | |
| self._portal: BlockingPortal | None = None | |
| self._portal_lock: Lock = Lock() | |
| def has_session_credentials(self) -> bool: | |
| return self._session_key is not None | |
| async def prompt(self, value: str) -> RuntimeSnapshot: | |
| return await to_thread.run_sync(self.prompt_sync, value) | |
| async def _prompt(self, value: str) -> RuntimeSnapshot: | |
| if self._closed: | |
| raise RuntimeError("Space runtime session is closed") | |
| prompt = value.strip() | |
| if not prompt: | |
| raise ValueError("prompt must not be empty") | |
| if self._prompt_lock is None: | |
| self._prompt_lock = anyio.Lock() | |
| async with self._prompt_lock: | |
| await self._start() | |
| _ = await self.manager.emit( | |
| "generate", GenerateData(prompt=prompt), source="space" | |
| ) | |
| await self.manager.wait_for_idle(timeout=120.0 if self.provider == "bonsai" else 30.0) | |
| return self.snapshot() | |
| def prompt_sync(self, value: str) -> RuntimeSnapshot: | |
| return self._blocking_portal().call(self._prompt, value) | |
| async def close(self) -> None: | |
| await to_thread.run_sync(self.close_sync) | |
| async def _close(self) -> None: | |
| if not self._closed: | |
| try: | |
| await self.manager.wait_for_idle() | |
| await self.manager.close() | |
| finally: | |
| self._session_key = None | |
| self._closed = True | |
| def close_sync(self) -> None: | |
| with self._portal_lock: | |
| portal = self._portal | |
| if portal is None: | |
| stack = ExitStack() | |
| self._portal_stack = stack | |
| portal = stack.enter_context(start_blocking_portal()) | |
| self._portal = portal | |
| stack = self._portal_stack | |
| if stack is None: | |
| raise RuntimeError("Space runtime portal stack is unavailable") | |
| try: | |
| portal.call(self._close) | |
| finally: | |
| stack.close() | |
| self._portal = None | |
| self._portal_stack = None | |
| def snapshot(self) -> RuntimeSnapshot: | |
| secrets = (self._session_key,) if self._session_key is not None else () | |
| events = tuple( | |
| _event_record(event, secrets=secrets) for event in self.manager.get_events() | |
| ) | |
| notifications = tuple( | |
| _payload_record(record, kind="error") | |
| for record in events | |
| if record.get("name") == "error" | |
| ) | |
| plans = tuple( | |
| _payload_record(record) | |
| for record in events | |
| if record.get("name") == "space_plan" | |
| ) | |
| agent_statuses = _latest_by(events, "voice_status", "voice") | |
| workflows = tuple( | |
| _payload_record(record) | |
| for record in events | |
| if record.get("name") | |
| in { | |
| "workflow_started", | |
| "workflow_phase_start", | |
| "workflow_phase_complete", | |
| "workflow_completed", | |
| } | |
| ) | |
| artifacts = tuple( | |
| _payload_record(record) | |
| for record in events | |
| if record.get("name") == "space_artifact" | |
| ) | |
| plugins: tuple[JsonRecord, ...] = tuple( | |
| { | |
| "name": plugin.name, | |
| "enabled": plugin.enabled, | |
| "capabilities": list(plugin.capabilities), | |
| } | |
| for plugin in self.manager.get_plugins() | |
| ) | |
| mcp: tuple[JsonRecord, ...] = ( | |
| {"name": "mcp", "status": "fixture", "servers": []}, | |
| ) | |
| return RuntimeSnapshot( | |
| provider=self.provider, | |
| model=self.model, | |
| mode=self.mode, | |
| closed=self._closed, | |
| events=events, | |
| notifications=notifications, | |
| plans=plans, | |
| agent_statuses=agent_statuses, | |
| workflows=workflows, | |
| mcp=mcp, | |
| plugins=plugins, | |
| artifacts=artifacts, | |
| ) | |
| async def _start(self) -> None: | |
| if self._started: | |
| return | |
| for plugin in self._plugins: | |
| await self.manager.enable_plugin(plugin) | |
| self._started = True | |
| def _blocking_portal(self) -> BlockingPortal: | |
| with self._portal_lock: | |
| if self._portal is None: | |
| stack = ExitStack() | |
| self._portal_stack = stack | |
| self._portal = stack.enter_context(start_blocking_portal()) | |
| return self._portal | |
| def create_runtime_session( | |
| selection: ProviderSelection, | |
| builder: RuntimeBuilder | None = None, | |
| ) -> SpaceRuntimeSession: | |
| selected_builder = _fixture_builder if builder is None else builder | |
| manager, plugins = selected_builder(selection) | |
| runtime_mode = os.getenv("KATETO_SPACE_MODE", "live") | |
| model = "fixture/demo-model" | |
| if selection.provider == "fixture": | |
| runtime_mode = "fixture" | |
| if builder is None and selection.provider != "fixture": | |
| config = SpaceProviderConfig.from_env() | |
| manager, plugins = _live_builder(selection, config) | |
| model = ( | |
| config.byok_model | |
| if selection.provider == "byok" | |
| else config.bonsai_model or "unconfigured" | |
| ) | |
| elif runtime_mode not in {"fixture", "live"}: | |
| raise ValueError("KATETO_SPACE_MODE must be fixture or live") | |
| return SpaceRuntimeSession( | |
| selection, manager, plugins, mode=runtime_mode, model=model | |
| ) | |
| def _event_record( | |
| envelope: EventEnvelope[BaseModel], *, secrets: tuple[str, ...] = () | |
| ) -> JsonRecord: | |
| data: JsonValue = envelope.data.model_dump(mode="json") | |
| record: JsonRecord = { | |
| "name": envelope.name, | |
| "source": envelope.source, | |
| "timestamp": envelope.timestamp.isoformat(), | |
| "data": _redact(data, secrets=secrets), | |
| } | |
| return record | |
| def _redact(value: JsonValue, *, secrets: tuple[str, ...] = ()) -> JsonValue: | |
| if isinstance(value, dict): | |
| return {key: _redact(item, secrets=secrets) for key, item in value.items()} | |
| if isinstance(value, list): | |
| return [_redact(item, secrets=secrets) for item in value] | |
| if isinstance(value, str): | |
| for secret in secrets: | |
| if secret: | |
| value = value.replace(secret, "<redacted>") | |
| return value | |
| def _latest_by( | |
| events: tuple[JsonRecord, ...], event_name: str, key: str | |
| ) -> tuple[JsonRecord, ...]: | |
| latest: dict[str, JsonRecord] = {} | |
| for record in events: | |
| if record.get("name") != event_name: | |
| continue | |
| data = record.get("data") | |
| if isinstance(data, dict): | |
| value = data.get(key) | |
| if isinstance(value, str): | |
| latest[value] = record | |
| return tuple(latest.values()) | |
| def _payload_record(record: JsonRecord, *, kind: str | None = None) -> JsonRecord: | |
| payload = record.get("data") | |
| flattened: JsonRecord = {"name": record["name"]} | |
| if isinstance(payload, dict): | |
| flattened.update(payload) | |
| if kind is not None: | |
| flattened["kind"] = kind | |
| return flattened | |
| async def close_runtime_session(session: SpaceRuntimeSession | None) -> None: | |
| if session is not None: | |
| await session.close() | |