KatetoSpace / space /runtime.py
Chaos
fix(space): allow Bonsai startup to finish
7797763
Raw
History Blame Contribute Delete
19.8 kB
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, "<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()