Spaces:
Running on Zero
Running on Zero
File size: 19,820 Bytes
7573c9c 2e52a00 7573c9c 35a216d 2578b07 7573c9c 2578b07 7573c9c 35a216d 7573c9c 35a216d 7573c9c 2e52a00 2578b07 7573c9c 2e52a00 7573c9c 35a216d 7573c9c 2578b07 7573c9c 2578b07 7573c9c 2578b07 7573c9c 2578b07 7573c9c 2e52a00 2578b07 7573c9c 2e52a00 2578b07 2e52a00 2578b07 2e52a00 cf46599 2e52a00 2578b07 2e52a00 7573c9c 2578b07 7573c9c 2578b07 7573c9c 2578b07 7573c9c 2578b07 7573c9c 2578b07 7573c9c 2578b07 7573c9c 7797763 7573c9c cf46599 7573c9c 2e52a00 cf46599 7573c9c | 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 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 | 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()
|