This repository has no description
Something went wrong. Try again.
9.6 kB · 272 lines
Python
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273"""Execution-scoped RPC authority and bounded output, not persistent memory."""from __future__ import annotations
import asyncioimport contextvarsimport ioimport uuidfrom dataclasses import dataclass, fieldfrom typing import Any, Literal, get_args
from .wire import Wire, fields
OUTPUT_LIMIT = 65_536OUTPUT_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."""
@dataclassclass 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)