This repository has no description
Something went wrong. Try again.
14 kB · 287 lines
Python
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288from __future__ import annotations
import asyncioimport jsonimport sqlite3import tempfileimport unittestfrom pathlib import Path
from klbr import behavior, hooksfrom klbr.behavior_runner import EphemeralRunner, TestReportfrom peer import ROOT, Peer, behavior_copy, output
class BehaviorIntegrationTests(unittest.TestCase): def setUp(self): self.temp_dir = tempfile.TemporaryDirectory(prefix="klbr-integ-") self.tmp = Path(self.temp_dir.name) self.releases_dir = self.tmp / "releases" self.db = sqlite3.connect(":memory:") self.db.execute("PRAGMA foreign_keys=ON") self.db.executescript((ROOT / "klbr-runtime/migrations/001_runtime.sql").read_text()) self.db.execute("INSERT INTO sessions(id) VALUES('test_session')") self.db.commit()
def tearDown(self): self.db.close() self.temp_dir.cleanup()
def test_full_lifecycle_progression(self): # 1. Draft draft_a = behavior.draft(ROOT / "behavior") self.assertIsInstance(draft_a, behavior.Draft) fp_a = draft_a.fingerprint()
# 2. test() runs EphemeralRunner, returns ValidatedCandidate candidate_a = draft_a.test(timeout=5.0) self.assertIsInstance(candidate_a, behavior.ValidatedCandidate) self.assertEqual(candidate_a.fingerprint, fp_a) self.assertIsInstance(candidate_a.test_report, TestReport) self.assertEqual(candidate_a.test_report.fingerprint, fp_a) self.assertIn("attention.operator", candidate_a.test_report.results) self.assertIn("delivery.operator", candidate_a.test_report.results)
# 3. publish() produces immutable PublishedRelease release_a = candidate_a.publish(self.releases_dir) self.assertIsInstance(release_a, behavior.PublishedRelease) self.assertEqual(release_a.id, fp_a) self.assertTrue(release_a.path.exists()) self.assertTrue(release_a.verify())
# 4. activate() updates SQLite sessions table and logs event active_1 = behavior.activate("test_session", self.db, expected_epoch=0, release=release_a) self.assertIsInstance(active_1, behavior.ActiveEpoch) self.assertEqual(active_1.epoch, 1) self.assertEqual(active_1.release.id, fp_a)
# Verify database state after activation session_row = self.db.execute("SELECT hook_revision, hook_epoch FROM sessions WHERE id='test_session'").fetchone() self.assertEqual(session_row, (fp_a, 1)) events = self.db.execute("SELECT seq, kind, payload FROM events WHERE session_id='test_session' ORDER BY seq").fetchall() self.assertEqual(len(events), 1) self.assertEqual(events[0][1], "hooks.activated") self.assertEqual(json.loads(events[0][2]), {"revision": fp_a, "epoch": 1})
# 5. Create, test, publish, and activate second release draft_b = behavior.draft(ROOT / "behavior") draft_b.write("klbr_hooks/custom.py", "# custom behavior extension\n") draft_b.write("manifest.json", json.dumps({ "version": 1, "hooks": { "attention.plan": "klbr_hooks.defaults:attention", "delivery.review": "klbr_hooks.defaults:delivery", } })) candidate_b = draft_b.test(timeout=5.0) release_b = candidate_b.publish(self.releases_dir) self.assertNotEqual(release_b.id, release_a.id)
active_2 = behavior.activate("test_session", self.db, expected_epoch=1, release=release_b) self.assertEqual(active_2.epoch, 2) self.assertEqual(active_2.release.id, release_b.id)
session_row = self.db.execute("SELECT hook_revision, hook_epoch FROM sessions WHERE id='test_session'").fetchone() self.assertEqual(session_row, (release_b.id, 2))
# 6. rollback() restores release_a at incremented epoch 3 with rollback flag in event active_3 = behavior.rollback("test_session", self.db, expected_epoch=2, release=release_a) self.assertEqual(active_3.epoch, 3) self.assertEqual(active_3.release.id, release_a.id)
session_row = self.db.execute("SELECT hook_revision, hook_epoch FROM sessions WHERE id='test_session'").fetchone() self.assertEqual(session_row, (release_a.id, 3))
events = self.db.execute("SELECT seq, kind, payload FROM events WHERE session_id='test_session' ORDER BY seq").fetchall() self.assertEqual(len(events), 3) self.assertEqual(json.loads(events[2][2]), {"revision": release_a.id, "epoch": 3, "rollback": True})
def test_cas_conflict_on_stale_epoch(self): draft = behavior.draft(ROOT / "behavior") release = behavior.publish(draft, self.releases_dir)
# Expected epoch 0 succeeds active = behavior.activate("test_session", self.db, expected_epoch=0, release=release) self.assertEqual(active.epoch, 1)
# Same expected epoch 0 now fails with CASConflictError with self.assertRaises(behavior.CASConflictError) as ctx: behavior.activate("test_session", self.db, expected_epoch=0, release=release) self.assertIn("stale epoch: expected 0", str(ctx.exception))
# Session remains untouched at epoch 1 session_row = self.db.execute("SELECT hook_revision, hook_epoch FROM sessions WHERE id='test_session'").fetchone() self.assertEqual(session_row, (release.id, 1))
# Attempting rollback with stale epoch also fails with self.assertRaises(behavior.CASConflictError): behavior.rollback("test_session", self.db, expected_epoch=0, release=release)
def test_cas_conflict_on_nonexistent_session(self): draft = behavior.draft(ROOT / "behavior") release = behavior.publish(draft, self.releases_dir) with self.assertRaises(behavior.CASConflictError): behavior.activate("missing_session", self.db, expected_epoch=0, release=release)
def test_ephemeral_runner_catches_infinite_loop(self): draft = behavior.draft(ROOT / "behavior") draft.write("klbr_hooks/defaults.py", """from klbr.hooks import AttentionDecision, AttentionRequest, DeliveryDecision, DeliveryRequest, HookContext
def attention(request: AttentionRequest, context: HookContext) -> AttentionDecision: while True: pass
def delivery(request: DeliveryRequest, context: HookContext) -> DeliveryDecision: return DeliveryDecision.allow()""") with self.assertRaises(behavior.ValidationError) as ctx: EphemeralRunner.run(draft, timeout=0.5) self.assertIn("timed out", str(ctx.exception).lower())
def test_ephemeral_runner_catches_blocking_sleep(self): draft = behavior.draft(ROOT / "behavior") draft.write("klbr_hooks/defaults.py", """import timefrom klbr.hooks import AttentionDecision, AttentionRequest, DeliveryDecision, DeliveryRequest, HookContext
def attention(request: AttentionRequest, context: HookContext) -> AttentionDecision: time.sleep(30.0) return AttentionDecision.wake("operator")
def delivery(request: DeliveryRequest, context: HookContext) -> DeliveryDecision: return DeliveryDecision.allow()""") with self.assertRaises(behavior.ValidationError) as ctx: EphemeralRunner.run(draft, timeout=0.5) self.assertIn("timed out", str(ctx.exception).lower())
def test_ephemeral_runner_catches_exception_in_hook(self): draft = behavior.draft(ROOT / "behavior") draft.write("klbr_hooks/defaults.py", """from klbr.hooks import AttentionDecision, AttentionRequest, DeliveryDecision, DeliveryRequest, HookContext
def attention(request: AttentionRequest, context: HookContext) -> AttentionDecision: raise ValueError("intentional hook crash")
def delivery(request: DeliveryRequest, context: HookContext) -> DeliveryDecision: return DeliveryDecision.allow()""") with self.assertRaises(behavior.ValidationError) as ctx: EphemeralRunner.run(draft, timeout=1.0) self.assertIn("intentional hook crash", str(ctx.exception))
def test_ephemeral_runner_catches_syntax_error(self): draft = behavior.draft(ROOT / "behavior") draft.write("klbr_hooks/defaults.py", "def broken syntax !!!") with self.assertRaises(behavior.ValidationError) as ctx: EphemeralRunner.run(draft, timeout=1.0) self.assertTrue("syntax" in str(ctx.exception).lower() or "worker startup failed" in str(ctx.exception).lower())
def test_ephemeral_runner_catches_invalid_return_type(self): draft = behavior.draft(ROOT / "behavior") draft.write("klbr_hooks/defaults.py", """from klbr.hooks import AttentionDecision, AttentionRequest, DeliveryDecision, DeliveryRequest, HookContext
def attention(request: AttentionRequest, context: HookContext) -> AttentionDecision: return {"not": "a typed AttentionDecision"}
def delivery(request: DeliveryRequest, context: HookContext) -> DeliveryDecision: return DeliveryDecision.allow()""") with self.assertRaises(behavior.ValidationError) as ctx: EphemeralRunner.run(draft, timeout=1.0) self.assertIn("AttentionDecision", str(ctx.exception))
def test_ephemeral_runner_catches_invalid_delivery_decision(self): draft = behavior.draft(ROOT / "behavior") draft.write("klbr_hooks/defaults.py", """from klbr.hooks import AttentionDecision, AttentionRequest, DeliveryDecision, DeliveryRequest, HookContext
def attention(request: AttentionRequest, context: HookContext) -> AttentionDecision: return AttentionDecision.wake("operator")
def delivery(request: DeliveryRequest, context: HookContext) -> DeliveryDecision: # hold with empty reason is invalid return DeliveryDecision("hold", "")""") with self.assertRaises(behavior.ValidationError) as ctx: EphemeralRunner.run(draft, timeout=1.0) self.assertIn("invalid delivery decision", str(ctx.exception).lower())
class PinningLawIntegrationTests(unittest.IsolatedAsyncioTestCase): """Pinning law: workbench cell retains its pinned revision/epoch across concurrent activation."""
async def test_workbench_execution_retains_pinning_across_activation(self): with tempfile.TemporaryDirectory(prefix="klbr-pinning-") as tmp: tmp_path = Path(tmp) releases_dir = tmp_path / "releases" db = sqlite3.connect(":memory:") db.execute("PRAGMA foreign_keys=ON") db.executescript((ROOT / "klbr-runtime/migrations/001_runtime.sql").read_text()) db.execute("INSERT INTO sessions(id) VALUES('pin_session')") db.commit()
# Release 1: default behavior (delivery returns allow) d1 = behavior.draft(ROOT / "behavior") rel1 = behavior.publish(d1, releases_dir) active1 = behavior.activate("pin_session", db, 0, rel1) self.assertEqual(active1.epoch, 1)
# Start an execution pinned to rel1 ev_start = db.execute( "INSERT INTO events(session_id, kind, payload, created_ms) VALUES('pin_session', 'execution.started', '{}', 0)" ).lastrowid db.execute( "INSERT INTO executions(id, session_id, generation, hook_revision, state, started_event) VALUES('exec_1', 'pin_session', 'gen_1', ?, 'running', ?)", (rel1.id, ev_start) ) db.commit()
# Release 2: custom delivery behavior d2 = behavior.draft(ROOT / "behavior") d2.write("klbr_hooks/defaults.py", """from klbr.hooks import AttentionDecision, AttentionRequest, DeliveryDecision, DeliveryRequest, HookContext
def attention(request: AttentionRequest, context: HookContext) -> AttentionDecision: return AttentionDecision.wake("operator")
def delivery(request: DeliveryRequest, context: HookContext) -> DeliveryDecision: return DeliveryDecision.hold("policy v2 holds everything")""") rel2 = behavior.publish(d2, releases_dir)
# Concurrent activation happens mid-execution: session advances to epoch 2, revision rel2 active2 = behavior.activate("pin_session", db, 1, rel2) self.assertEqual(active2.epoch, 2) self.assertEqual(active2.release.id, rel2.id)
# PINNING LAW: The in-flight execution 'exec_1' is still pinned to rel1.id! pinned_rev = db.execute("SELECT hook_revision FROM executions WHERE id='exec_1'").fetchone()[0] self.assertEqual(pinned_rev, rel1.id) self.assertNotEqual(pinned_rev, rel2.id)
# Complete execution 1 ev_end = db.execute( "INSERT INTO events(session_id, kind, payload, created_ms) VALUES('pin_session', 'execution.completed', '{}', 1)" ).lastrowid db.execute("UPDATE executions SET state='ok', completed_event=? WHERE id='exec_1'", (ev_end,)) db.commit()
# Next execution starts with the new session revision rel2 ev_start_2 = db.execute( "INSERT INTO events(session_id, kind, payload, created_ms) VALUES('pin_session', 'execution.started', '{}', 2)" ).lastrowid current_session_rev = db.execute("SELECT hook_revision FROM sessions WHERE id='pin_session'").fetchone()[0] self.assertEqual(current_session_rev, rel2.id)
db.execute( "INSERT INTO executions(id, session_id, generation, hook_revision, state, started_event) VALUES('exec_2', 'pin_session', 'gen_2', ?, 'running', ?)", (current_session_rev, ev_start_2) ) db.commit()
pinned_rev_2 = db.execute("SELECT hook_revision FROM executions WHERE id='exec_2'").fetchone()[0] self.assertEqual(pinned_rev_2, rel2.id) db.close()