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 @final @dataclass(frozen=True, slots=True) 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) @final 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 @override 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() @property 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, "") 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()