from __future__ import annotations import asyncio import json import sqlite3 import tempfile import unittest from pathlib import Path from klbr import behavior, hooks from klbr.behavior_runner import EphemeralRunner, TestReport from 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 time from 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()