ive harnessed the harness
Something went wrong. Try again.
10 kB · 213 lines
Python
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214"""Entrypoint for both supervised roles. Python >=3.11, Unix control socket."""from __future__ import annotations
# Source-tree launch works with Python -I -B; no inherited PYTHONPATH required.if __package__ in {None, ""}: import sys from pathlib import Path sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) __package__ = "klbr_runtime"
import argparseimport asyncioimport contextlibimport jsonimport sysimport tracebackfrom pathlib import Pathfrom typing import Any
from klbr.hooks import HookContextfrom . import skill_indexfrom .context import Bridge, RoutedTextfrom .kernel import Kernelfrom .policy import Policyfrom .wire import ( MAX_CODE, BindingToken, ProtocolError, RELAY_PROTOCOL_VERSION, VERSION, Wire, fields, identifier, measured_runtime_identity, text,)
async def serve(args: argparse.Namespace) -> None: # Read the host's skill payload before anything else: the file lives beside the # control socket in a directory the host owns, and the host drops that directory # once this process has connected. Leaving the read until a cell asks would race # the very handshake that makes the directory disposable. if args.skills: skill_index.install(args.skills) if args.behavior: behavior_path = str(Path(args.behavior).resolve()) if behavior_path not in sys.path: sys.path.insert(0, behavior_path) policy = Policy(Path(args.behavior), args.revision) if args.role == "hooks" else None reader, writer = await asyncio.open_unix_connection(args.socket) # The --target value is the host-issued launch binding: a binding token # this process echoes so the host can refuse a foreign or stale worker. # It is NOT measured identity; the runtime identity below is measured # on this machine by the endpoint itself. binding = BindingToken.parse(json.loads(args.target)) if args.target else None wire = Wire(reader, writer, args.session, args.generation, target=binding.as_stamp() if binding else None) bridge, kernel, dispatcher = Bridge(wire), None, None if args.role == "workbench": # The tree vocabulary is bound before the first cell runs, so ordinary work does not # begin with an import line the model has to remember. The package stays importable for # anything it does not pre-bind. from klbr import workspace as tree
kernel = Kernel(bridge, prelude=tree.prelude()) elif args.role == "remote": # An auxiliary remote namespace: a persistent target-side Python heap that is # deliberately NOT the workbench. It shares no globals with a local cell, and # the dispatcher is the only place that knows what an operation name means. from klbr_runtime.remote_dispatch import RemoteDispatchError, RemoteDispatcher
dispatcher = RemoteDispatcher(workspace=args.workspace) hello: dict[str, Any] = {"kind": "hello", "token": args.token, "role": args.role, "revision": args.revision, "slots": sorted(policy.handlers) if policy else [], # Echoed, not assumed: the host records what this generation # actually loaded rather than what the roots hold now. "skills": skill_index.revision() if skill_index.loaded() else None} if binding is not None: runtime = measured_runtime_identity() if dispatcher is not None: # Capabilities belong INSIDE the measured runtime identity, because that is # the endpoint's own probe of this machine. The host reads them from # identity.runtime.capabilities to decide admission, so an advertised-but- # unusable capability is refused before a call rather than at use. runtime["capabilities"] = list(dispatcher.capabilities) hello["identity"] = { "protocol_version": VERSION, "relay_protocol_version": RELAY_PROTOCOL_VERSION, "target": {"stamp": binding.as_stamp()}, "runtime": runtime, } await wire.send(hello) original_out, original_err = sys.stdout, sys.stderr sys.stdout, sys.stderr = RoutedText("stdout", original_out), RoutedText("stderr", original_err) active: asyncio.Task | None = None active_id: str | None = None
async def run(message: dict) -> None: try: if message["kind"] == "execute": assert kernel is not None outcome = await kernel.execute( message["id"], message["code"], purpose=message.get("purpose", "interactive"), ) elif message["kind"] == "remote_dispatch": assert dispatcher is not None # The runner frames, bounds and identifies; it never interprets the # operation name. Dispatch owns the whole vocabulary. try: value = await dispatcher.dispatch( message["operation"], message["payload"], operation_id=message["id"], ) outcome = {"kind": "remote", "value": value} except RemoteDispatchError as error: outcome = {"kind": "remote_error", "code": error.code, "message": error.message} else: assert policy is not None ctx = HookContext(message["id"], args.session, args.revision) value = await policy.invoke(message["slot"], message["input"], ctx) outcome = {"kind": "hook", "value": value} except asyncio.CancelledError: outcome = {"kind": "failed", "error": "hook invocation interrupted"} except BaseException as error: outcome = {"kind": "failed", "error": "".join(traceback.format_exception(error))[-16_384:]} await wire.send({"kind": "completed", "id": message["id"], "outcome": outcome})
try: while True: message = await wire.receive() kind = message["kind"] if kind == "host_reply": bridge.reply(message) elif kind == "interrupt": fields(message, {"kind", "id"}) if active and not active.done() and message["id"] == active_id: active.cancel() elif kind == "shutdown": fields(message, {"kind"}) break elif kind in {"execute", "invoke", "remote_dispatch"}: if kind == "remote_dispatch": fields(message, {"kind", "id", "operation", "payload"}) if dispatcher is None: raise ProtocolError("only a remote worker can dispatch operations") text(message["operation"], 128) elif kind == "execute": fields(message, {"kind", "id", "code"}, {"purpose"}) if kernel is None: raise ProtocolError("hook worker cannot execute cells") text(message["code"], MAX_CODE) purpose = message.get("purpose", "interactive") if purpose not in {"interactive", "maintenance", "eval"}: raise ProtocolError(f"invalid execution purpose: {purpose}") else: fields(message, {"kind", "id", "slot", "input"}) if policy is None: raise ProtocolError("workbench cannot run required policy hooks") text(message["slot"], 128) identifier(message["id"]) if active and not active.done(): raise ProtocolError("concurrent foreground operation rejected") if active: active.result() # Propagate previous transport errors. active_id = message["id"] active = asyncio.create_task(run(message)) else: raise ProtocolError(f"unexpected host message: {kind}") finally: if active and not active.done(): active.cancel() # Host enforces a hard process deadline; do not depend on cooperative # cancellation when arbitrary Python catches CancelledError. sys.stdout, sys.stderr = original_out, original_err writer.close() with contextlib.suppress(ConnectionError): await writer.wait_closed()
def main() -> None: parser = argparse.ArgumentParser() parser.add_argument("--socket", required=True) parser.add_argument("--role", choices=["workbench", "hooks", "remote"], required=True) parser.add_argument("--session", required=True) parser.add_argument("--generation", required=True) parser.add_argument("--token", required=True) # Host-issued launch binding (a TargetStamp JSON). Echoed in Hello as a # binding token; verified identity always comes from a measured probe. parser.add_argument("--target") parser.add_argument("--behavior") parser.add_argument("--workspace") parser.add_argument("--revision") # Content-identity payload of the skill bodies this worker may offer. Written by # the host in the same directory as the control socket; absent means no skills. parser.add_argument("--skills") args = parser.parse_args() if args.role == "hooks" and (not args.behavior or not args.revision): parser.error("hook workers need a behavior snapshot and revision") sys.dont_write_bytecode = True try: asyncio.run(serve(args)) except (ProtocolError, ConnectionError, asyncio.IncompleteReadError) as error: print(f"klbr worker disconnected: {error}", file=sys.stderr) raise SystemExit(2) from error
if __name__ == "__main__": main()