ive harnessed the harness
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460"""Ephemeral dry-run test runner, synthetic effect harness, and CAS switch for klbr behavior.
Runs isolated validation against candidate behavior drafts before publication or activation."""from __future__ import annotations
import asyncioimport concurrent.futuresimport contextlibimport jsonimport mathimport osfrom pathlib import Pathimport shutilimport sqlite3import sysimport tempfileimport timeimport uuidfrom dataclasses import dataclassfrom typing import Any
from klbr.hooks import AttentionDecision, DeliveryDecisionfrom klbr_runtime.wire import ProtocolError, Wire
try: from klbr.behavior import ( ActiveEpoch, CASConflictError, Draft, PublishedRelease, ValidationError, )except ImportError: class ValidationError(ValueError): # type: ignore[no-redef] pass
class CASConflictError(RuntimeError): # type: ignore[no-redef] pass
class Draft: # type: ignore[no-redef] pass
class PublishedRelease: # type: ignore[no-redef] pass
class ActiveEpoch: # type: ignore[no-redef] pass
@dataclass(frozen=True, slots=True)class TestReport: fingerprint: str slots: tuple[str, ...] results: dict[str, Any] duration_seconds: float
def _validate_attention_decision(value: Any, expected_priority: str | None = None) -> AttentionDecision: if not isinstance(value, dict): raise ValidationError(f"expected dict for AttentionDecision, got {type(value).__name__}") action = value.get("action") if action == "wake": priority = value.get("priority") if priority not in {"operator", "directed", "ambient"}: raise ValidationError(f"invalid priority for wake: {priority!r}") if set(value.keys()) != {"action", "priority"}: raise ValidationError(f"unexpected keys in wake decision: {set(value.keys())}") if expected_priority is not None and priority != expected_priority: raise ValidationError(f"expected priority {expected_priority!r}, got {priority!r}") return AttentionDecision.wake(priority) elif action == "defer": seconds = value.get("seconds") if not isinstance(seconds, (int, float)) or not math.isfinite(seconds) or not (0 < seconds <= 600): raise ValidationError(f"invalid defer seconds: {seconds!r}") if set(value.keys()) != {"action", "seconds"}: raise ValidationError(f"unexpected keys in defer decision: {set(value.keys())}") return AttentionDecision.defer(float(seconds)) elif action == "ignore": reason = value.get("reason") if not isinstance(reason, str) or not reason.strip(): raise ValidationError(f"invalid ignore reason: {reason!r}") if set(value.keys()) != {"action", "reason"}: raise ValidationError(f"unexpected keys in ignore decision: {set(value.keys())}") return AttentionDecision.ignore(reason) else: raise ValidationError(f"unknown AttentionDecision action: {action!r}")
def _validate_delivery_decision(value: Any) -> DeliveryDecision: if not isinstance(value, dict): raise ValidationError(f"expected dict for DeliveryDecision, got {type(value).__name__}") action = value.get("action") if action == "allow": if set(value.keys()) != {"action"}: raise ValidationError(f"unexpected keys in allow decision: {set(value.keys())}") return DeliveryDecision.allow() elif action == "hold": reason = value.get("reason") if not isinstance(reason, str) or not reason.strip(): raise ValidationError(f"invalid hold reason: {reason!r}") if set(value.keys()) != {"action", "reason"}: raise ValidationError(f"unexpected keys in hold decision: {set(value.keys())}") return DeliveryDecision.hold(reason) else: raise ValidationError(f"unknown DeliveryDecision action: {action!r}")
def _resolve_runner_script() -> Path: # Try locating __main__.py relative to this file or cwd candidate = Path(__file__).resolve().parent.parent / "klbr_runtime" / "__main__.py" if candidate.exists(): return candidate cwd_candidate = Path.cwd() / "python/klbr-runtime/src/klbr_runtime/__main__.py" if cwd_candidate.exists(): return cwd_candidate.resolve() raise RuntimeError(f"cannot locate klbr_runtime entrypoint (checked {candidate})")
class EphemeralRunner: """Ephemeral hook worker test runner and synthetic effect harness."""
def __init__(self, timeout: float = 2.0, startup_timeout: float = 5.0) -> None: self.default_timeout = timeout self.startup_timeout = startup_timeout
@classmethod def run( cls, draft: Any, timeout: float = 2.0, startup_timeout: float = 5.0, ) -> TestReport: runner = cls(timeout=timeout, startup_timeout=startup_timeout) return runner.execute(draft, timeout=timeout, startup_timeout=startup_timeout)
def execute( self, draft: Any, timeout: float | None = None, startup_timeout: float | None = None, ) -> TestReport: effective_timeout = self.default_timeout if timeout is None else timeout effective_startup = self.startup_timeout if startup_timeout is None else startup_timeout try: asyncio.get_running_loop() except RuntimeError: return asyncio.run(self.run_async(draft, timeout=effective_timeout, startup_timeout=effective_startup)) # If an event loop is already running in this thread, execute in a thread pool with concurrent.futures.ThreadPoolExecutor(max_workers=1) as pool: return pool.submit( lambda: asyncio.run(self.run_async(draft, timeout=effective_timeout, startup_timeout=effective_startup)) ).result()
async def run_async( self, draft: Any, timeout: float = 2.0, startup_timeout: float = 5.0, ) -> TestReport: start_time = time.monotonic() fp = draft.fingerprint() manifest = draft.manifest()
# Export draft to temporary directory if in-memory or copy draft_temp_dir = tempfile.TemporaryDirectory(prefix="klbr-draft-") socket_temp_dir = tempfile.TemporaryDirectory(prefix="klbr-sock-") server = None process = None writer = None drain_tasks = []
try: draft_path = Path(draft_temp_dir.name) / "behavior" draft.export(draft_path)
socket_path = str(Path(socket_temp_dir.name) / "control.sock") loop = asyncio.get_running_loop() accepted: asyncio.Future[tuple[asyncio.StreamReader, asyncio.StreamWriter]] = loop.create_future()
def on_accept(reader: asyncio.StreamReader, w: asyncio.StreamWriter) -> None: if not accepted.done(): accepted.set_result((reader, w)) else: w.close()
server = await asyncio.start_unix_server(on_accept, socket_path)
session = "test" generation = str(uuid.uuid4()) token = str(uuid.uuid4())
py_bin = os.environ.get("KLBR_PYTHON", sys.executable) runner_script = _resolve_runner_script()
args = [ py_bin, "-I", "-B", "-u", str(runner_script), "--role", "hooks", "--socket", socket_path, "--session", session, "--generation", generation, "--token", token, "--behavior", str(draft_path), "--revision", fp, ]
env = { k: os.environ[k] for k in ("PATH", "HOME", "LANG", "LC_ALL", "TMPDIR", "LD_LIBRARY_PATH", "DYLD_LIBRARY_PATH") if k in os.environ }
process = await asyncio.create_subprocess_exec( *args, env=env, stdin=asyncio.subprocess.DEVNULL, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, )
stderr_buf = bytearray()
async def drain_stream(stream: asyncio.StreamReader | None, buf: bytearray | None = None) -> None: if stream is None: return while chunk := await stream.read(4096): if buf is not None: buf += chunk if len(buf) > 16384: del buf[:-16384]
drain_tasks = [ asyncio.create_task(drain_stream(process.stdout)), asyncio.create_task(drain_stream(process.stderr, stderr_buf)), ]
# Handshake with dedicated startup timeout (independent of invocation timeout) handshake_timeout = startup_timeout try: reader, writer = await asyncio.wait_for(accepted, timeout=handshake_timeout) except asyncio.TimeoutError as err: raise ValidationError( f"worker startup timed out after {handshake_timeout}s; stderr: {stderr_buf.decode(errors='replace')}" ) from err except BaseException as err: raise ValidationError( f"worker startup failed: {err}; stderr: {stderr_buf.decode(errors='replace')}" ) from err
wire = Wire(reader, writer, session, generation) try: hello = await asyncio.wait_for(wire.receive(), timeout=handshake_timeout) except asyncio.TimeoutError as err: raise ValidationError(f"worker hello handshake timed out after {handshake_timeout}s") from err except (ProtocolError, Exception) as err: raise ValidationError(f"worker hello handshake error: {err}") from err
if hello.get("kind") != "hello": raise ValidationError(f"expected hello message, got {hello.get('kind')}") if hello.get("token") != token: raise ValidationError("worker token mismatch in hello") if hello.get("role") != "hooks": raise ValidationError(f"expected role 'hooks', got {hello.get('role')!r}") if hello.get("revision") != fp: raise ValidationError(f"worker revision mismatch: {hello.get('revision')!r} != {fp!r}")
slots = hello.get("slots", []) if "attention.plan" not in slots or "delivery.review" not in slots: raise ValidationError(f"worker missing required slots: {slots}")
# Synthetic effect harness async def invoke(slot: str, payload: dict) -> dict: op_id = str(uuid.uuid4()) await wire.send({"kind": "invoke", "id": op_id, "slot": slot, "input": payload}) reply = await asyncio.wait_for(wire.receive(), timeout=timeout) if reply.get("kind") != "completed" or reply.get("id") != op_id: raise ValidationError(f"unexpected worker reply for {slot}: {reply}") outcome = reply.get("outcome", {}) if outcome.get("kind") != "hook": err_msg = outcome.get("error", "hook invocation failed") raise ValidationError(f"hook invocation {slot!r} failed: {err_msg}") return outcome.get("value")
results: dict[str, Any] = {} hook_checks = ( ( "attention.plan", {"event_id": 1, "source": "operator", "received_ms": 0}, "attention.plan(operator)", _validate_attention_decision, "attention.operator", ), ( "attention.plan", {"event_id": 2, "source": "discord.mention", "received_ms": 0}, "attention.plan(discord.mention)", _validate_attention_decision, "attention.mention", ), ( "attention.plan", {"event_id": 3, "source": "discord.ambient", "received_ms": 0}, "attention.plan(discord.ambient)", _validate_attention_decision, "attention.ambient", ), ( "delivery.review", {"effect_id": "synthetic-test", "destination": "operator", "text": "synthetic test"}, "delivery.review", _validate_delivery_decision, "delivery.operator", ), ) for slot, payload, label, validator, key in hook_checks: try: val = await invoke(slot, payload) except asyncio.TimeoutError as err: raise ValidationError(f"{label} timed out after {timeout}s") from err except ValidationError: raise except Exception as err: raise ValidationError(f"{label} error: {err}") from err results[key] = validator(val)
duration = round(time.monotonic() - start_time, 4) return TestReport( fingerprint=fp, slots=tuple(sorted(manifest.keys())), results=results, duration_seconds=duration, )
finally: if writer is not None: writer.close() with contextlib.suppress(Exception): await writer.wait_closed() if server is not None: server.close() with contextlib.suppress(Exception): await server.wait_closed() if process is not None and process.returncode is None: process.kill() with contextlib.suppress(Exception): await process.wait() for t in drain_tasks: t.cancel() with contextlib.suppress(asyncio.CancelledError): await t draft_temp_dir.cleanup() socket_temp_dir.cleanup()
def activate( session_id: str, db: sqlite3.Connection | Path | str, expected_epoch: int, release: PublishedRelease,) -> ActiveEpoch: """Activate a release for session_id via SQLite CAS transaction.""" return _cas_update( session_id=session_id, db=db, expected_epoch=expected_epoch, release=release, is_rollback=False, )
def rollback( session_id: str, db: sqlite3.Connection | Path | str, expected_epoch: int, target_release: PublishedRelease,) -> ActiveEpoch: """Roll back session_id to target_release via SQLite CAS transaction.""" return _cas_update( session_id=session_id, db=db, expected_epoch=expected_epoch, release=target_release, is_rollback=True, )
def _cas_update( session_id: str, db: sqlite3.Connection | Path | str, expected_epoch: int, release: PublishedRelease, *, is_rollback: bool = False,) -> ActiveEpoch: """Execute BEGIN IMMEDIATE, release insert, session CAS, and hooks.activated event.""" release.verify() if isinstance(db, (Path, str)): conn = sqlite3.connect(str(db)) close_conn = True elif isinstance(db, sqlite3.Connection): conn = db close_conn = False else: raise TypeError(f"expected sqlite3.Connection, Path, or str, got {type(db).__name__}")
old_isolation = conn.isolation_level try: conn.isolation_level = None # Autocommit mode for explicit transaction control conn.execute("PRAGMA foreign_keys = ON") conn.execute("BEGIN IMMEDIATE") try: conn.execute( "INSERT INTO releases(id, path) VALUES(?, ?) ON CONFLICT DO NOTHING", (release.id, str(release.path)), ) cur = conn.execute( "UPDATE sessions SET hook_revision = ?, hook_epoch = hook_epoch + 1 WHERE id = ? AND hook_epoch = ?", (release.id, session_id, expected_epoch), ) if cur.rowcount != 1: raise CASConflictError(f"stale epoch: expected {expected_epoch}")
new_epoch = expected_epoch + 1 now_ms = int(time.time() * 1000) payload_dict: dict[str, Any] = {"revision": release.id, "epoch": new_epoch} if is_rollback: payload_dict["rollback"] = True payload = json.dumps(payload_dict, separators=(",", ":")) conn.execute( "INSERT INTO events(session_id, kind, payload, created_ms) VALUES(?, ?, ?, ?)", (session_id, "hooks.activated", payload, now_ms), ) conn.execute("COMMIT") except Exception: with contextlib.suppress(sqlite3.OperationalError): conn.execute("ROLLBACK") raise finally: conn.isolation_level = old_isolation if close_conn: conn.close()
return ActiveEpoch( session_id=session_id, epoch=new_epoch, release=release, db=db, )