Skip to content

Commit efb7af1

Browse files
fix: wire plugin lifecycle hooks across surfaces
Fixes #18
1 parent 7dc259d commit efb7af1

8 files changed

Lines changed: 406 additions & 38 deletions

File tree

‎Makefile‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,7 @@ clean:
4545

4646
# Syntax check
4747
lint:
48-
$(VENV_PYTHON) -m compileall -q main.py server.py tools.py memory.py event_store.py task_engine.py learning_engine.py learning_benchmark.py migration.py scheduler.py skill_registry.py browser_manager.py instruction_loader.py subagents.py capability_tokens.py tool_registry.py dynamic_tools.py
48+
$(VENV_PYTHON) -m compileall -q main.py server.py tools.py memory.py event_store.py task_engine.py learning_engine.py learning_benchmark.py migration.py scheduler.py skill_registry.py browser_manager.py instruction_loader.py subagents.py capability_tokens.py tool_registry.py dynamic_tools.py plugin_runtime.py
4949
@echo "Python syntax OK."
5050
@echo "All files pass syntax check."
5151

@@ -56,7 +56,7 @@ test:
5656
# Quick verification
5757
check:
5858
@echo "Checking Python syntax..."
59-
@$(VENV_PYTHON) -m py_compile main.py server.py tools.py memory.py event_store.py task_engine.py learning_engine.py learning_benchmark.py migration.py scheduler.py skill_registry.py browser_manager.py instruction_loader.py subagents.py capability_tokens.py tool_registry.py dynamic_tools.py
59+
@$(VENV_PYTHON) -m py_compile main.py server.py tools.py memory.py event_store.py task_engine.py learning_engine.py learning_benchmark.py migration.py scheduler.py skill_registry.py browser_manager.py instruction_loader.py subagents.py capability_tokens.py tool_registry.py dynamic_tools.py plugin_runtime.py
6060
@echo " Python modules: OK"
6161
@echo "Checking git tools..."
6262
@$(VENV_PYTHON) -c "from tools import AVAILABLE_TOOLS; git = [k for k in AVAILABLE_TOOLS if k.startswith('git_')]; print(f' {len(git)} git tools, {len(AVAILABLE_TOOLS)} total tools')"

‎README.md‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -551,6 +551,23 @@ def register():
551551

552552
Available hooks: `on_startup`, `on_turn_start`, `on_turn_end`, `on_tool_execute`.
553553

554+
The shared runtime loads `plugins/*.py` once per execution surface (`cli` and
555+
`web`) in deterministic filename order. Hook failures are isolated, recorded
556+
as `plugin.hook_failed` events, and never fail the user turn. Hook arguments
557+
are keyword arguments and are bounded/redacted before a plugin sees them:
558+
559+
| Hook | Required kwargs | Additional context |
560+
|------|-----------------|--------------------|
561+
| `on_startup` | — | `agent`, `surface`, `user_id`, `workspace_id` |
562+
| `on_turn_start` | `user_input` | `surface`, `user_id`, `workspace_id`, `session_id`, `profile` |
563+
| `on_turn_end` | `reply`, `success` | `error` plus the turn context |
564+
| `on_tool_execute` | `action`, `args`, `result`, `success` | `error`, `surface`, and scope context |
565+
566+
`on_tool_execute` fires once for every attempted action, including unknown,
567+
unauthorized, malformed, and approval-denied actions. `turn_logger.py` writes
568+
to `kyrozen_turns.log`; set `KYROZEN_TURN_LOG` to choose another path for a
569+
service or test.
570+
554571
See `plugins/turn_logger.py` for a working example.
555572

556573
---

‎main.py‎

Lines changed: 85 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,7 @@ def _terminal_supports_unicode() -> bool:
102102
from subagents import AgentProfile, SubAgentManager
103103
from capability_tokens import issue_capability_token
104104
from dynamic_tools import SAFE_BUILTINS, validate_tool_source
105+
from plugin_runtime import get_plugin_runtime
105106
from tools import (AVAILABLE_TOOLS, set_workspace_root as _set_tools_workspace_root,
106107
resolve_capabilities, tool_capability)
107108
from providers import (
@@ -1875,6 +1876,23 @@ def _system_prompt(tools_list: str) -> str:
18751876
_learning_notices: list[str] = []
18761877

18771878

1879+
def _record_plugin_event(event_type: str, payload: dict[str, Any]) -> None:
1880+
"""Persist plugin runtime diagnostics without allowing them to break work."""
1881+
try:
1882+
memory_bank.store.append_event(
1883+
event_type, payload, user_id=memory_bank.user_id,
1884+
workspace_id=memory_bank.workspace_id, session_id=memory_bank.session_id,
1885+
)
1886+
except Exception:
1887+
pass
1888+
1889+
1890+
def _plugin_runtime_for_surface(surface: str | None = None):
1891+
return get_plugin_runtime(
1892+
surface or _EXECUTION_SURFACE, event_recorder=_record_plugin_event,
1893+
)
1894+
1895+
18781896
def _subagent_action_arguments(action: str, args: Any) -> Any:
18791897
"""Normalise common structured action arguments without widening access."""
18801898
if not isinstance(args, dict):
@@ -1914,12 +1932,17 @@ def _run_subagent_tool(profile: AgentProfile, action: str, args: Any, tools: set
19141932
canonical = TOOL_ALIASES.get(str(action).strip(), str(action).strip())
19151933
receipt_id = f"subagent_receipt_{uuid.uuid4().hex}"
19161934
if canonical not in tools:
1935+
_notify_tool_execute(
1936+
canonical, args,
1937+
f"Error: tool '{canonical}' is not authorized for the '{profile.name}' profile.",
1938+
)
19171939
return {
19181940
"receipt_id": receipt_id, "action": canonical, "args": str(args)[:500],
19191941
"result": f"Error: tool '{canonical}' is not authorized for the '{profile.name}' profile.",
19201942
"success": False, "authorized": False,
19211943
}
19221944
if not isinstance(args, (str, dict)):
1945+
_notify_tool_execute(canonical, args, "Error: Action args must be a plain string or supported object.")
19231946
return {
19241947
"receipt_id": receipt_id, "action": canonical, "args": str(args)[:500],
19251948
"result": "Error: Action args must be a plain string or supported object.",
@@ -1978,6 +2001,10 @@ def _run_subagent_llm(profile: AgentProfile, task: str, context: list[dict[str,
19782001
if unknown or malformed:
19792002
executed_steps += 1
19802003
rejected_action = unknown or "malformed"
2004+
_notify_tool_execute(
2005+
str(rejected_action), "",
2006+
"Error: malformed or unknown Action rejected; no tool was executed.",
2007+
)
19812008
tool_records.append({
19822009
"receipt_id": f"subagent_receipt_{uuid.uuid4().hex}",
19832010
"action": str(rejected_action), "args": "",
@@ -2561,6 +2588,22 @@ def _add(data: dict) -> None:
25612588
return calls
25622589

25632590

2591+
def _notify_tool_execute(action: str, args: Any, result: Any) -> None:
2592+
"""Emit exactly one isolated plugin hook for every attempted tool call."""
2593+
result_text = str(result)
2594+
try:
2595+
_plugin_runtime_for_surface().tool_execute(
2596+
action=action, args=args, result=result_text,
2597+
success=not _is_tool_error(result_text),
2598+
error=None if not _is_tool_error(result_text) else result_text,
2599+
user_id=memory_bank.user_id, workspace_id=memory_bank.workspace_id,
2600+
session_id=memory_bank.session_id,
2601+
)
2602+
except Exception:
2603+
# A plugin is an observer and cannot change tool semantics.
2604+
pass
2605+
2606+
25642607
def _run_tool(action: str, args: str) -> str:
25652608
# Map aliases
25662609
action = TOOL_ALIASES.get(action, action)
@@ -2593,15 +2636,21 @@ def _run_tool(action: str, args: str) -> str:
25932636
pass
25942637
fn = AVAILABLE_TOOLS.get(action)
25952638
if not fn:
2596-
return f"Error: unknown tool '{action}'"
2639+
result = f"Error: unknown tool '{action}'"
2640+
_notify_tool_execute(action, args, result)
2641+
return result
25972642
required_capability = tool_capability(action)
25982643
if not _execution_capability_token.allows(required_capability):
2599-
return f"Error: tool '{action}' requires capability '{required_capability}'"
2644+
result = f"Error: tool '{action}' requires capability '{required_capability}'"
2645+
_notify_tool_execute(action, args, result)
2646+
return result
26002647
if not _confirm_tool_action(action, str(args)):
2601-
return (
2648+
result = (
26022649
f"Error: {action} requires confirmation. "
26032650
"Approve it interactively or set KYROZEN_APPROVAL_MODE=never for an explicitly automated CLI."
26042651
)
2652+
_notify_tool_execute(action, args, result)
2653+
return result
26052654
start = time.time()
26062655
try:
26072656
result = str(fn(args))
@@ -2611,6 +2660,7 @@ def _run_tool(action: str, args: str) -> str:
26112660
result = f"Error: {e}"
26122661
elapsed = time.time() - start
26132662
_track_tool_performance(action, result, elapsed)
2663+
_notify_tool_execute(action, args, result)
26142664
return result
26152665

26162666

@@ -3115,8 +3165,8 @@ def _classify_complexity(user_input: str) -> str:
31153165
return "medium"
31163166

31173167

3118-
def _chat_turn(user_input: str, clear_tasks: bool = False, profile: str | None = None,
3119-
memory_context: dict[str, Any] | None = None) -> str:
3168+
def _chat_turn_impl(user_input: str, clear_tasks: bool = False, profile: str | None = None,
3169+
memory_context: dict[str, Any] | None = None) -> str:
31203170
"""One user turn: build context, get LLM reply, execute tool calls
31213171
with automatic retries and failure memory."""
31223172

@@ -3189,6 +3239,9 @@ def _chat_turn(user_input: str, clear_tasks: bool = False, profile: str | None =
31893239
if not unknown_action:
31903240
break
31913241
_unknown_retries += 1
3242+
_notify_tool_execute(
3243+
unknown_action, "", f"Error: unknown tool '{unknown_action}'",
3244+
)
31923245
msg = (
31933246
f"System: Action '{unknown_action}' is not recognized. "
31943247
"You must use one of the following actions: "
@@ -3315,6 +3368,7 @@ def _chat_turn(user_input: str, clear_tasks: bool = False, profile: str | None =
33153368
"Example: `\"args\": \"python3 process_logs.py\"`\n"
33163369
"Do NOT use `\"args\": {\"cmd\": ...}`.\n"
33173370
)
3371+
_notify_tool_execute(action, args, result)
33183372
else:
33193373
if not isinstance(args, str):
33203374
args = str(args)
@@ -3581,6 +3635,7 @@ def _chat_turn(user_input: str, clear_tasks: bool = False, profile: str | None =
35813635
"Example: `\"args\": \"python3 process_logs.py\"`\n"
35823636
"Do NOT use `\"args\": {\"cmd\": ...}`.\n"
35833637
)
3638+
_notify_tool_execute(action, args, result2)
35843639
else:
35853640
if not isinstance(args, str):
35863641
args = str(args)
@@ -3697,6 +3752,30 @@ def _chat_turn(user_input: str, clear_tasks: bool = False, profile: str | None =
36973752
)
36983753

36993754

3755+
def _chat_turn(user_input: str, clear_tasks: bool = False, profile: str | None = None,
3756+
memory_context: dict[str, Any] | None = None) -> str:
3757+
"""Run one chat turn with failure-isolated plugin lifecycle hooks."""
3758+
runtime = _plugin_runtime_for_surface()
3759+
context = {
3760+
"user_id": memory_bank.user_id,
3761+
"workspace_id": memory_bank.workspace_id,
3762+
"session_id": memory_bank.session_id,
3763+
"profile": profile or _agent_profile_mode,
3764+
}
3765+
runtime.turn_start(user_input=user_input, **context)
3766+
try:
3767+
reply = _chat_turn_impl(
3768+
user_input, clear_tasks=clear_tasks, profile=profile,
3769+
memory_context=memory_context,
3770+
)
3771+
except Exception as exc:
3772+
runtime.turn_end(reply="", success=False,
3773+
error=f"{type(exc).__name__}: {exc}", **context)
3774+
raise
3775+
runtime.turn_end(reply=reply, success=True, **context)
3776+
return reply
3777+
3778+
37003779
def _split_reply(text: str) -> tuple[str, str]:
37013780
text = text.strip()
37023781
thought_match = re.search(r"^Thought:\s*(.*?)(?=\n(?:Action:|(?:\n|$)))", text, re.DOTALL | re.MULTILINE)
@@ -3965,6 +4044,7 @@ def main() -> None:
39654044
sys.exit(1)
39664045
# Set workspace root for sandbox
39674046
_set_workspace_root(os.getcwd())
4047+
_plugin_runtime_for_surface().load_once()
39684048
task_results = _run_recovered_tasks()
39694049
for item in task_results:
39704050
console.print(f"[{_SUCCESS}]Durable task {item['task']['id']}: {item['status']}[/{_SUCCESS}]")

‎plugin_runtime.py‎

Lines changed: 139 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,139 @@
1+
"""Shared, failure-isolated plugin lifecycle runtime."""
2+
3+
from __future__ import annotations
4+
5+
import importlib.util
6+
import re
7+
import threading
8+
from pathlib import Path
9+
from typing import Any, Callable
10+
11+
12+
_RUNTIMES: dict[str, "PluginRuntime"] = {}
13+
_RUNTIMES_LOCK = threading.RLock()
14+
15+
_SECRET_PATTERNS = (
16+
re.compile(r"(?i)((?:api[_-]?key|password|secret|token)\s*[:=])\s*\S+"),
17+
re.compile(r"\bsk-[A-Za-z0-9_-]+"),
18+
)
19+
20+
21+
def _redact(value: Any, limit: int = 500) -> str:
22+
text = str(value or "")
23+
for pattern in _SECRET_PATTERNS:
24+
text = pattern.sub(lambda match: (match.group(1) + "<redacted>")
25+
if match.lastindex else "<redacted>", text)
26+
return text.replace("\n", " ")[:limit]
27+
28+
29+
class PluginRuntime:
30+
"""Load plugins once for one execution surface and isolate every hook."""
31+
32+
def __init__(self, surface: str, *, plugins_dir: Path | None = None,
33+
event_recorder: Callable[[str, dict[str, Any]], None] | None = None):
34+
self.surface = str(surface or "unknown")
35+
self.plugins_dir = (plugins_dir or Path(__file__).parent / "plugins").resolve()
36+
self.event_recorder = event_recorder
37+
self._loaded = False
38+
self._plugins: list[tuple[str, Any]] = []
39+
self._lock = threading.RLock()
40+
41+
@property
42+
def loaded(self) -> bool:
43+
return self._loaded
44+
45+
@property
46+
def plugin_names(self) -> tuple[str, ...]:
47+
return tuple(name for name, _plugin in self._plugins)
48+
49+
def _record(self, event_type: str, **payload: Any) -> None:
50+
if not self.event_recorder:
51+
return
52+
safe_payload = {key: _redact(value, 300) if isinstance(value, str) else value
53+
for key, value in payload.items()}
54+
try:
55+
self.event_recorder(event_type, safe_payload)
56+
except Exception:
57+
pass
58+
59+
def load_once(self) -> tuple[str, ...]:
60+
with self._lock:
61+
if self._loaded:
62+
return self.plugin_names
63+
# Set the flag before importing anything: a broken plugin must not
64+
# cause another request to retry imports or duplicate registrations.
65+
self._loaded = True
66+
if not self.plugins_dir.is_dir():
67+
self._record("plugin.loaded", plugin="<none>", status="no plugin directory")
68+
return self.plugin_names
69+
for path in sorted(self.plugins_dir.glob("*.py")):
70+
if path.name.startswith("_"):
71+
continue
72+
module_name = f"openkyrozen_plugin_{self.surface}_{path.stem}"
73+
try:
74+
spec = importlib.util.spec_from_file_location(module_name, path)
75+
if spec is None or spec.loader is None:
76+
raise ImportError("could not create module spec")
77+
module = importlib.util.module_from_spec(spec)
78+
spec.loader.exec_module(module)
79+
register = getattr(module, "register", None)
80+
if not callable(register):
81+
self._record("plugin.load_failed", plugin=path.stem,
82+
error="register() is missing")
83+
continue
84+
plugin = register()
85+
if plugin is None:
86+
self._record("plugin.load_failed", plugin=path.stem,
87+
error="register() returned no plugin")
88+
continue
89+
self._plugins.append((path.stem, plugin))
90+
self._record("plugin.loaded", plugin=path.stem, status="loaded")
91+
except Exception as exc:
92+
self._record("plugin.load_failed", plugin=path.stem,
93+
error=f"{type(exc).__name__}: {exc}")
94+
return self.plugin_names
95+
96+
def trigger(self, hook_name: str, **kwargs: Any) -> None:
97+
self.load_once()
98+
for plugin_name, plugin in tuple(self._plugins):
99+
hook = getattr(plugin, hook_name, None)
100+
if not callable(hook):
101+
continue
102+
try:
103+
hook(**kwargs)
104+
except Exception as exc:
105+
self._record(
106+
"plugin.hook_failed", plugin=plugin_name, hook=hook_name,
107+
error=f"{type(exc).__name__}: {exc}",
108+
)
109+
110+
def turn_start(self, *, user_input: str, **context: Any) -> None:
111+
self.trigger("on_turn_start", user_input=_redact(user_input, 1000),
112+
surface=self.surface, **context)
113+
114+
def turn_end(self, *, reply: str, success: bool, error: str | None = None,
115+
**context: Any) -> None:
116+
self.trigger("on_turn_end", reply=_redact(reply, 1000), success=bool(success),
117+
error=_redact(error, 300) if error else None,
118+
surface=self.surface, **context)
119+
120+
def tool_execute(self, *, action: str, args: Any, result: Any, success: bool,
121+
error: str | None = None, **context: Any) -> None:
122+
self.trigger(
123+
"on_tool_execute", action=_redact(action, 100), args=_redact(args, 500),
124+
result=_redact(result, 1000), success=bool(success),
125+
error=_redact(error, 300) if error else None,
126+
surface=self.surface, **context,
127+
)
128+
129+
130+
def get_plugin_runtime(surface: str, *, plugins_dir: Path | None = None,
131+
event_recorder: Callable[[str, dict[str, Any]], None] | None = None) -> PluginRuntime:
132+
"""Return the one runtime instance for a surface in this process."""
133+
key = str(surface or "unknown")
134+
with _RUNTIMES_LOCK:
135+
runtime = _RUNTIMES.get(key)
136+
if runtime is None:
137+
runtime = PluginRuntime(key, plugins_dir=plugins_dir, event_recorder=event_recorder)
138+
_RUNTIMES[key] = runtime
139+
return runtime

‎plugins/turn_logger.py‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,13 +4,14 @@
44
"""
55

66
import time
7+
import os
78
from pathlib import Path
89

910
class TurnLogger:
1011
"""Logs every conversation turn with timestamp."""
1112

1213
def __init__(self):
13-
self._log_path = Path("kyrozen_turns.log")
14+
self._log_path = Path(os.environ.get("KYROZEN_TURN_LOG", "kyrozen_turns.log"))
1415

1516
def on_turn_start(self, user_input: str, **kwargs):
1617
ts = time.strftime("%Y-%m-%d %H:%M:%S")

‎pyproject.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,7 @@ Repository = "https://github.com/EvanProgramming/OpenKyrozen"
4848
Issues = "https://github.com/EvanProgramming/OpenKyrozen/issues"
4949

5050
[tool.setuptools]
51-
py-modules = ["main", "memory", "providers", "server", "tools", "event_store", "learning_engine", "learning_benchmark", "migration", "task_engine", "scheduler", "skill_registry", "browser_manager", "instruction_loader", "subagents", "capability_tokens", "tool_registry", "dynamic_tools"]
51+
py-modules = ["main", "memory", "providers", "server", "tools", "event_store", "learning_engine", "learning_benchmark", "migration", "task_engine", "scheduler", "skill_registry", "browser_manager", "instruction_loader", "subagents", "capability_tokens", "tool_registry", "dynamic_tools", "plugin_runtime"]
5252
include-package-data = true
5353

5454
[tool.setuptools.packages.find]

0 commit comments

Comments
 (0)