"""Execution-scoped RPC authority and bounded output, not persistent memory.""" from __future__ import annotations import asyncio import contextvars import io import uuid from dataclasses import dataclass, field from typing import Any, Literal, get_args from .wire import Wire, fields OUTPUT_LIMIT = 65_536 OUTPUT_RECORD_LIMIT = 256 PURPOSE_INTERACTIVE = "interactive" PURPOSE_MAINTENANCE = "maintenance" PURPOSE_EVAL = "eval" _PURPOSES = (PURPOSE_INTERACTIVE, PURPOSE_MAINTENANCE, PURPOSE_EVAL) def _validate_purpose(purpose: str) -> str: if purpose not in _PURPOSES: raise ValueError(f"invalid execution purpose: {purpose}") return purpose class ScopeClosed(RuntimeError): """The originating foreground execution no longer owns host-effect authority.""" class HostError(RuntimeError): def __init__(self, code: str, message: str): super().__init__(f"{code}: {message}") self.code = code class TurnYield(BaseException): """Interpreter control outcome. Not a pickled/resumable continuation.""" @dataclass class ExecutionScope: id: str bridge: Bridge open: bool = True terminal: str | None = None output: list[dict[str, str]] = field(default_factory=list) stored_bytes: int = 0 truncated_bytes: int = 0 purpose: str = PURPOSE_INTERACTIVE def __post_init__(self) -> None: self.purpose = _validate_purpose(self.purpose) def write(self, stream: str, value: str) -> None: raw = value.encode("utf-8", errors="replace") room = max(0, OUTPUT_LIMIT - self.stored_bytes) kept = raw[:room].decode("utf-8", errors="ignore") count = len(kept.encode("utf-8")) if kept: if self.output and self.output[-1]["stream"] == stream: self.output[-1]["text"] += kept elif len(self.output) < OUTPUT_RECORD_LIMIT: self.output.append({"stream": stream, "text": kept}) else: count = 0 self.stored_bytes += count self.truncated_bytes += len(raw) - count async def request(self, request: dict) -> Any: if not self.open or self.terminal is not None: raise ScopeClosed("host operations require a live, non-yielded foreground execution") return await self.bridge.request(self, request) current: contextvars.ContextVar[ExecutionScope | None] = contextvars.ContextVar("klbr_execution", default=None) def scope() -> ExecutionScope: value = current.get() if value is None or not value.open or value.terminal is not None: raise ScopeClosed("no active workbench execution; hooks return proposals instead of calling effects") return value def current_purpose() -> str: value = current.get() return value.purpose if value is not None else PURPOSE_INTERACTIVE async def originate_maintenance( purpose: str, reason: str, evidence_refs: list[str] | None = None, impact: str | None = None, ) -> dict: s = scope() payload: dict[str, Any] = { "method": "maintenance.originate", "purpose": purpose, "reason": reason, "evidence_refs": evidence_refs or [], } if impact is not None: payload["impact"] = impact return await s.request(payload) # What the host records against a work item, mirroring its MaintenanceAssessment variants # exactly. One definition, in the type: the tuple below is derived from it, so the check and the # annotation cannot drift apart. # # This is a *disposition* set, not a verdict about a diagnosis: `klbr.maintenance` judges a claim # ("validated", "needs-evidence") and this decides the item's fate. They overlap in words and mean # different things, so a verdict submitted here is refused with a message that says so rather than # being coerced into the nearest-looking tag. Disposition = Literal[ "no-change", "needs-evidence", "deferred", "candidate-failed", "activated", ] ASSESSMENT_OUTCOMES: tuple[str, ...] = get_args(Disposition) async def assess_maintenance( work_item_id: str, outcome: Disposition, reason: str | None = None, missing_refs: list[str] | None = None, revisit_condition: str | None = None, revisit_deadline_ms: int | None = None, trigger_ms: int | None = None, candidate_revision: str | None = None, check_receipts: list[str] | None = None, ) -> dict: missing_refs = missing_refs or [] check_receipts = check_receipts or [] if outcome not in ASSESSMENT_OUTCOMES: raise ValueError( f"{outcome!r} is not a disposition for a work item; the host records one of " f"{', '.join(ASSESSMENT_OUTCOMES)}. A verdict about a diagnosis " f"(`klbr.maintenance.Verdict`) is a different set at a different level and has to be " f"mapped onto one of these explicitly." ) if outcome in ("no-change", "deferred", "candidate-failed") and not (reason or "").strip(): raise ValueError(f"{outcome} assessment requires a non-empty reason") if outcome == "needs-evidence" and not missing_refs and revisit_deadline_ms is None: raise ValueError( "needs-evidence assessment requires the actual missing refs " "or a bounded revisit_deadline_ms" ) if outcome == "activated" and not (candidate_revision or "").strip(): raise ValueError("activated assessment requires the candidate_revision it activated") s = scope() payload: dict[str, Any] = { "method": "maintenance.assess", "work_item_id": work_item_id, "outcome": outcome, "reason": reason, "missing_refs": missing_refs, "revisit_condition": revisit_condition, "revisit_deadline_ms": revisit_deadline_ms, "trigger_ms": trigger_ms, "candidate_revision": candidate_revision, "check_receipts": check_receipts, } return await s.request(payload) async def resolve_evidence(refs: list[str]) -> dict: """Resolve evidence reference IDs to their actual content. Returns {"resolved": [...], "unavailable": [...]} Each resolved entry has: ref, kind, and kind-specific content (e.g. event). Each unavailable entry has: ref, reason. """ s = scope() return await s.request({"method": "evidence.resolve", "refs": refs}) async def workspace_info() -> dict: """Get workspace path and session identity. Returns {"path": str, "session_id": str} """ s = scope() return await s.request({"method": "workspace.info"}) class Bridge: def __init__(self, wire: Wire): self.wire = wire self.pending: dict[str, tuple[str, asyncio.Future]] = {} async def request(self, execution: ExecutionScope, request: dict) -> Any: request_id = str(uuid.uuid4()) future = asyncio.get_running_loop().create_future() self.pending[request_id] = (execution.id, future) try: await self.wire.send({"kind": "host_request", "id": request_id, "execution_id": execution.id, "request": request}) return await future finally: self.pending.pop(request_id, None) if not future.done(): future.cancel() def reply(self, message: dict) -> None: fields(message, {"kind", "id", "execution_id", "reply"}) entry = self.pending.get(message["id"]) if entry is None: # A cancellation may precede its already-in-flight reply. return execution_id, future = entry if execution_id != message["execution_id"] or future.done(): return reply = message["reply"] if reply.get("kind") == "ok": fields(reply, {"kind", "value"}) future.set_result(reply["value"]) elif reply.get("kind") == "error": fields(reply, {"kind", "code", "message"}) future.set_exception(HostError(reply["code"], reply["message"])) else: raise ValueError("invalid host reply") def close_execution(self, execution: ExecutionScope) -> None: execution.open = False for request_id, (execution_id, future) in list(self.pending.items()): if execution_id == execution.id: if not future.done(): future.cancel() self.pending.pop(request_id, None) class RoutedText(io.TextIOBase): """Python text is cell output. Native fd writes remain raw diagnostics. Late background output is diagnostic, tagged with its originating execution; it must not become the output of whichever cell happens to run next. """ def __init__(self, stream: str, original: io.TextIOBase): self.stream, self.original = stream, original @property def encoding(self) -> str: return "utf-8" def writable(self) -> bool: return True def fileno(self) -> int: return self.original.fileno() def flush(self) -> None: self.original.flush() def write(self, value: str) -> int: if not isinstance(value, str): raise TypeError("text output requires str") execution = current.get() if execution is not None and execution.open: execution.write(self.stream, value) else: origin = execution.id if execution else "unassigned" # The host continuously drains this pipe, retaining only a bounded tail. self.original.write(f"[background origin={origin}] {value[:4096]}") self.original.flush() return len(value)