From 7c553b097395a93faf77939c244dcb7c8c10713c Mon Sep 17 00:00:00 2001 From: "firat.sertgoz" Date: Tue, 17 Feb 2026 09:40:05 +0400 Subject: [PATCH] feat: Add lifecycle hooks system with 6 interception points (#18) * feat: Add lifecycle hooks system with 6 interception points Implement extensible hook infrastructure for intercepting and transforming agent operations at well-defined points in the lifecycle: - BeforeInbound: intercept/modify/reject incoming user messages - BeforeToolCall: intercept/modify/reject tool executions (chat + job) - BeforeOutbound: intercept/modify/suppress outgoing responses - TransformResponse: transform final response before completing a turn - OnSessionStart: fire-and-forget notification on new session creation - OnSessionEnd: fire-and-forget notification on session pruning Hooks execute in priority order with modification chaining, reject short-circuits, configurable failure modes (FailOpen/FailClosed), and per-hook timeouts. Empty registry is zero-cost (all hooks pass through immediately). Co-Authored-By: Claude Opus 4.6 * fix: enforce hook fail-closed semantics * Merge upstream/main into feat/hooks-system-clean Resolve merge conflicts: - FEATURE_PARITY.md: Keep both upstream cron/routines status and hooks status - src/error.rs: Keep both Hook and Orchestrator/Worker error variants Co-Authored-By: Claude Opus 4.6 * fix: resolve CI test failures in pairing store and wizard - Fix pairing store truncate bug: record_failed_approve used .truncate(true) which wiped the file before reading, causing rate limiting to never accumulate past 1 attempt. Changed to .truncate(false) to preserve existing data. - Fix wizard test: skip test_install_missing_bundled_channels when telegram WASM artifact specifically isn't available, not just when all channels are empty (whatsapp may exist without telegram). - Add workspace exclude for subcrate directories to prevent cargo from discovering them as workspace members during builds. Co-Authored-By: Claude Opus 4.6 * fix: address PR #18 review comments - Remove duplicate maybe_hydrate_thread call (rebase artifact) - Fix RwLock held across async hook execution in HookRegistry::run() - Add tracing::warn for silent JSON parse failures in hook modifications - Refactor execute_tool_inner to accept &WorkerDeps instead of 8 Arc params - Use real user_id from JobContext instead of job_id UUID in BeforeToolCall hook Co-Authored-By: Claude Opus 4.6 * fix: cargo fmt + remove tracked worktree breaking CI - Apply rustfmt formatting (method chain line breaks, match arm style) - Remove .claude/worktrees/ from git tracking (caused submodule error in CI) - Add .claude/worktrees/ to .gitignore Co-Authored-By: Claude Opus 4.6 --------- Co-authored-by: Firat Sertgoz Co-authored-by: Claude Opus 4.6 --- .gitignore | 3 + .sidecar/shells.json.lock | 0 .todos/command_usage.jsonl | 1 + .todos/config.json | 7 + .todos/config.json.lock | 0 .todos/db.lock | 0 .todos/issues.db | Bin 0 -> 299008 bytes .todos/issues.db-shm | Bin 0 -> 32768 bytes .todos/issues.db-wal | Bin 0 -> 337872 bytes Cargo.toml | 8 + FEATURE_PARITY.md | 22 +- src/agent/agent_loop.rs | 131 ++++++++- src/agent/scheduler.rs | 5 + src/agent/session_manager.rs | 61 +++- src/agent/worker.rs | 103 ++++--- src/error.rs | 3 + src/hooks/hook.rs | 199 +++++++++++++ src/hooks/mod.rs | 19 ++ src/hooks/registry.rs | 555 +++++++++++++++++++++++++++++++++++ src/lib.rs | 1 + src/main.rs | 7 +- src/setup/wizard.rs | 2 - 22 files changed, 1058 insertions(+), 69 deletions(-) create mode 100644 .sidecar/shells.json.lock create mode 100644 .todos/command_usage.jsonl create mode 100644 .todos/config.json create mode 100644 .todos/config.json.lock create mode 100644 .todos/db.lock create mode 100644 .todos/issues.db create mode 100644 .todos/issues.db-shm create mode 100644 .todos/issues.db-wal create mode 100644 src/hooks/hook.rs create mode 100644 src/hooks/mod.rs create mode 100644 src/hooks/registry.rs diff --git a/.gitignore b/.gitignore index e3ff7420..4febf1eb 100644 --- a/.gitignore +++ b/.gitignore @@ -4,6 +4,9 @@ .env.* !.env.example +# Claude Code worktrees +.claude/worktrees/ + target/ # WASM build artifacts (loaded from disk, not bundled) diff --git a/.sidecar/shells.json.lock b/.sidecar/shells.json.lock new file mode 100644 index 00000000..e69de29b diff --git a/.todos/command_usage.jsonl b/.todos/command_usage.jsonl new file mode 100644 index 00000000..448236da --- /dev/null +++ b/.todos/command_usage.jsonl @@ -0,0 +1 @@ +{"ts":"2026-02-12T13:55:06.864688+04:00","cmd":"init","session":"ses_b1bea6","ok":true,"dur_ms":85} diff --git a/.todos/config.json b/.todos/config.json new file mode 100644 index 00000000..58b374b4 --- /dev/null +++ b/.todos/config.json @@ -0,0 +1,7 @@ +{ + "pane_heights": [ + 0.3333333333333333, + 0.3333333333333333, + 0.3333333333333333 + ] +} \ No newline at end of file diff --git a/.todos/config.json.lock b/.todos/config.json.lock new file mode 100644 index 00000000..e69de29b diff --git a/.todos/db.lock b/.todos/db.lock new file mode 100644 index 00000000..e69de29b diff --git a/.todos/issues.db b/.todos/issues.db new file mode 100644 index 0000000000000000000000000000000000000000..675062ef84472ab0c45bcc2c8e9156461f9e91c9 GIT binary patch literal 299008 zcmeI*PjDO8eZX;me*h9BMcEd!VQjTa$u5bKk_b|NEGddYQ;@=#q+|+=CA*%@3RsaF zfmjH;)SpSxRFu@2HaT~C%q4AZophQ*bMLe>nYO3idg!H>9O6uxr0F!1_U->&U_qL+ z;xOQsktkx{?tA+_zxQ`na&Eq}s9TzR*Jy00mOPs}oa*XIeO{JRsZ^Kvy&!)3U$?~% zJ^lyrXFBq8yC1t!C;$9NuSo40{)K(cx#9mE{zd;!2G;s+^?yI}sBbL$qaGtO+xz40 z;hxi}Uv>Rc*OSx_#sBv5b!sFd%^yvt^jbxGWZtiedMWCrWjLRevazwD)htsf_& zY1d0xX?8aKe9}_yRJDgjMK{e&%}hA-^88A1u2hsub8jsc<%ARDiJUCH z^om?6-YUsARu-1#R&L907jI8F3Gqw3vzD)yHOrE5)mGtBWOhesyJ~xLi`~)SIQbr5kRxtIG@TtQIGlZeU9Gi_5#GUA?|i zT)4Ju-?7|$a$a63UM-5u=ZiOEl{Y8!O49W~-q8zPT>FY#zW7Q;x;@)2YVmt@Vxr6> z96VBddAHdV60&(^>)3Es8Xr$T>uzZhvBqE?ZrcpZof<$*-O${&;Rnuj51L+Ln=mc4 zVeK*Sn&=U+CgT-oXQ`O$s(Uwn>lq(+lf_j@iMjmw!x?F5{G)_S^tPz7`O?~zOP#xLx@Rswe<&lpTWA-#R$FDJ_1N(I z&@MA6v}ST>%jnNa$B(C>XS` z=+1}1Gr5m?nHOf~vVpDlvhF}C#w{_ouIXZ+si}2y-LTBqhrU>hV@V$;%Gv`FsH6zn zv%h^a{-_rpgYPudS~;<0RP=`RSaB*RFD#de*NQ8_2r%vPTl$73s=HC&{nWghSc?I; zZ@U3FW~dZmAN#kg-mG-=X!`j@Cob!1tzz80D>U3cGqGs+M?X#=wx881M$PMQ!R+(? z2VU1|Xd9|t(`#!@(^GIlMJwy}V$E!K!e*_kH7t=yyAwXKh@Hx8w>I^im4g3McI)-- ztTZx`e&()CtHzpX|L%#!$3C*V*aXU+@%GV)&fUQT*_?%F^QOAiG-`SsVw*{`-Z0i0 z;!^KM_7sU+e!M3mT^wmQsM$VB!T#R+{LAUAG&-98=Bwd+<1RMLu9&Ga(C(;e=~nYZ z;a0h#ndOFFx7Q3!jm{Sr)5oUQV*Rlv7)Ztqu@dq{w!yPh^PYEFy`dWo-Fobn)|%su z8UIAF5wgs7V^>wZqg9(bl7FSz5QAyk2B4P9THO-kW^+&WE)V8St>JI4l4%;+16_N# z(^=()rtWniyjid8`Py>TFq2l%6|Jf@-H+Q8HuRd&Xy@(MB0AwT_pI2YQ1^q-W(APT zPj+Rb%cJf3g4?bYcaj#{>rY=6-=rf$_D=h-{f#dK5I_I{1Q0*~0R#|0009IL=v0B6 zT`~6mof=;1g#ZEwAbEpN4<{0tg_000IagfB*srAb>#U3$Xw1 z{1{V11Q0*~0R#|0009ILKmY**`1k*`0R#|0009ILKmY**5I_I{1Ug@U_y3(AV`_*1 z0tg_000IagfB*srAb4KkE6Pp1$4Z zN~5FcZw*=MovNmr=B8%4-@E2ligTr+T$+1pu_(L8JOH)?M3mEzU8)y0xLKJK1iTB@~a#xIqF zlZ?7n^G>uL*PBWRPqS3>o_AWkp&Jd|dfYC@8#DfihF-Hwd11L!yjENZ@|t#UtE%47 zs@`n{cjI5FHnf_h_!!N}YPqb{EwxtmZnrsI->6r$4UvP^P)yA<+tM_&2fFrfr?bip zO|`U&qFVA55w_C8Qc?C}5Y%dZb!DZvTvF`Ro29v>8*VK(>yryjVZtRjuiM z+=LCiW;EJ)`?a{!5OC`)O+|LnjrL%&DFR=!S}t^2=hf>g#f59jc5{^{>VEX{@=Eb) zQLxMxZ+gvqLa*d=`K@EaS!sMc{jA&R8xM`fJwH%p?8D($Pl+Y%)M0AsMpJ8ck-P>G zG%i1qUaF-wnruIEthHJtDO{reH7vzkSCi?)RY{4t{QRMe^ll-Y(nYq9%==Z{(iC;m z5}itXR-92osc546R>uZBqsZdBiKOPQcv7%eA8Hxq2 zHGQY(xmWlDKs-inrA@`CwFzQFtE%>RsBAZ)g}Hcx&g$~QJFCSBzvL-7JVTT-CC75j zyUium?)@vO!O9`KPw$So__vDg#WHy^`;ZW`7xIO?j(F1X*&a@u3Fx3 z5f4MdXr7RQ!Eihnft@&bz*W6=PjryBJ;6)g%P_n9IJ87NsH+&C&J1Lvh2z_IJqOqC zbIymsGr5m?nb*sl%LcaI%eteX7%9X`VNDmKMNO@n>xN~+6983B*QPv)asWLvB zyKS{>Y;4#!;}2c&_(vy_cuGZYSdZOi*={zn^bJi^ccZ@hDRnoo&gIYdWu&F?cA<t72TgYhLfX|AE)F8rp`c z*Yw(2(;O3=P|?b|y$)(#a0Mr9iq%xZ5}CBi;S-D4sZ2(iKiV#0zOz#Bf65O2W9qIn zGLn8a>BOmOteN)jo>+YBBhlEL8riNCyXA{cVJZdxr%x4M$w;?n+uc@lY2wb!Vm?kd zc%=C9ZnG&QWb?}T!`n6M=oO+4v5zke|1wpxHzwcwAe>Hqa~qSbLr0uW^we-l^Ev%L zM`A(pk46KPv>Ebe+N42i`%YwOkCKPFkFmzz-GKD8>!qwTJDYy)w$^BDyoQr-=;atw z!U;Qx#V_$TnekD3J5)0{gmsl) zmDq*OLFbBGezGegT^?=MncXuLcak={#gmOb`{~9|&wr#o)BRGa=Wh=E$)W2*|337c zp&NrgAN;Gq<$?bi_=|zpb3e#^kW2M{qyI$qSK04n=QIDBS?{^j`|IwPdjGBW+r8gT zeK+;h)N=oK?7UK`$uryC!W~0xRKyVC{&{Bd^hY?k;GTSXYhc!S^6>P&)8*Ww?)TYP zJ93Zh)O1=mQu6%u#YOu8#`5*~xtm4V4Y3=Cs%lz_v#XA6n7swX9(&o-BcIDkV`J&( z?g~jQJF9y!q=!GnNBd)8!ZA_5e(~k+C1Ouj?1c6MhMnqV)VhJ4M9Aujc9WKZID74r3&tUcm98F9Andd1tpi{NYS(;U&X z#bY{m=$Y{Z7nZLSZ;5AuZBCZ2FZ+*G{9IH^M9;}7zv~r_7?gf)uXle}r|Q|u8R^iWRSNd^&#hm~N*9l&pY#R+^j9m2|7cCM{6n8@380@6 zIoS@Me|RFyBd-;*(s@x*Iw;9s>MB9|^N$o-3i4CJbM12cqZ2_6`Nnisnmm_&G8`16 zo9q5#IJZCg{cc@1Eu-;xx}~tl6+vV4vrxRNCj4B)TwOgScDka>Y^+SDu*4EiC6?%% z+P*wDCsEqM$&B>DxpuKlEYO?uo9l3L>*R^7bnaaG=~u%=bMr}w86NX$5P8<$eBDkH zuwS3uuXf#bo^)N<+`dMYlK=v>mTunTq=;LJ{M;*t5A0{LP7$rsJo?PPdzppD^L(GbKA2 zIg`5ob6qs2_g~LQvq#gL?aKGF8Q-bogfjE*AJ0gIqv?l1=KjdwZVLQQf6$V*pW>cp zXYC&e9%bg&kBP>6JpH6Q$lYyI?lUas=g5IKTFQtdhQ&kk>iH!#0!nc zA}6{MI@&H{_oRq`x9T}Fy4^Eav%D^zlqq`cuAvknA6~ovY3C99le{}#g{MVMdeyrZ zcPNh}cmw!7&3h*r+hW*nbswnJ&8DSiE`Lhe?g^AXn26i-2DHy?ef?AJ^T;LlCcWc+ z`e^+2ywsg0mvDFI5C8sqJ;Hf#54#>1L;9ap|k7ta+g70$hR z{FA3$qu_oVhT4_T>4~r_Y@~Cr?hFo}RY9%Nv&_r%x7U z3x%`i&z?QEhfC~#|DPWIX-fR!3jqWWKmY**5I_I{1Q0*~0R%cmV5ob{{=3KCa|rkS zKl}fV2`?2v009ILKmY**5I_I{1Q0;L5g2@}`_-4kzk0U+%PBK&y6^vgmKy#UJP06w z00IagfB*srAb?009ILKmY**5I_I{1Q19TVE>=& z!chbeKmY**5I_I{1Q0*~0R%cpfc<|bMU%Q9fB*srAbmu~;iNJMAbN1Q0*~0R#|0009ILKmY**I!eIa|EGt4nG(PFLI42-5I_I{1Q0*~0R#|00D%q@ z$aatPZ>W0B-T!}>8vd}uDoAw@KmY**5I_I{1Q0*~0R#{@xB`cJM_%o7HvXHndo|;s zv;X)0|NjoI3N=Rn0R#|0009ILKmY**5I~^A1%~YU9Zsc&|GUF0O7##x009ILKmY** z5I_I{1Q0k_0t4ypquF%#NMF@hGx_iT9jxF}TLch5009ILKmY**5I_KdgDJrN|6s6o0R#|0009ILKmY**5IC3u?EepD+^IDJ2q1s} z0tg_000IagfWW~LVE=!x;!SN4KmY**5I_I{1Q0*~0R#@F0Q>)g8Fy-p00IagfB*sr zAb_1~$S zP^;Awy0f1TclI;WGiOgs&zvaC$kXR9oH=u0`fS3Xq1=(wXzIJE!HI!?7%=*e_x(8g zx7o?e_cE7K|JVD|-oNg-)%{`D-LCV8>H`;Y|B_RNelqxYb`~$50y;#gSC%hN)Y+Q8N<{7QMnsb8jscIdEuSa;)H*XQ?h?Y zdHde;lQ%Qc<Ws#}_(Zd!)(S#hp1m4f@{7oKKTveM{i`uSSO?7Es~6Ch1n zP8MH&^<#2&LQUP!TvWg6^8EG1MLYb<*XQSM7N?x+tEy=!4|MHe+`~;QH#AkWofuS* zTYE7lQzgxwZu3G=CpgNT|H&Ps5E;7S(*>^ zjjCGH%<)`)?VYSN`$qanchHZ_$F;JuuA7$8c&u1v^yB%K4i!xmC)-^rdT7GZH#E~y zH|qJV#T!}ajW^QI-ghD%JuCY2o3ZFeQ^FSDH}Xg$kmpuQ*B6$>#Y>_&yABn-p_T0p z8>e*>RVPKU9@m?g++?k0>6WPe_Q`&0^c#tLt*JF0Xbr{G?)#zg?nir|mo>$|C(`Z&kt3iO@Lz=!;qDq9}8;U6~H7-Kc}h=XWa9 zz9eYN!DSPMi0i|KZdqE*o-3a8EM%q2qUaOtigue(yvA?sRJeOlLh+GCIVF1^^h$o~ z{1>uPKA(PeBJ5egJ;%$Q^U8K7sd(D1`&Z<4l6qKhgGq^D{QN!=nX>HhylRxysv^cz zRsQ15>&tErYG@l`l5aj5k4o5k{9d_TH=4bck*qZWc7`{ZAh4AOkKdD{KN@HW` zr{f{Gb5?lH#p1U9` zIDcFm+)VY^%fpR#v(hDzRew`fO4V3Xf+5hWwm8UM1@EJZ1VKc6nWm5~Y~>HER`i8`BM#q3X0 zEq4_^$G#)`NTk~8oLIz;L}zv9l7~j)o@ggdt?dJo$vH<)3o?%$v}^nJ=QC1%B)t~D zac4J{0VH|1hhVQ-aEbzY`6yUZeZ=vOUIMY%{>>Rrs=fFgCbymDN zpWo`cos}dh{nT)Juyd00?Shx1J3u;z!f~;6&2EoWVxdq?TtkR4SWFaF((o5qX*!wu zLzO$sMea5swXWA{N&VElR&A8m^#|?KBddEm?_BJAboD)KmY**5I_I{1Q0*~0R#|u;R3w>f8nHLhyVfzAbHQdicYX_{A3j2q1s}0tg_000IagfB*sr xbcVp;?(zORl@n^UdO|nNP0i%}e`f@iS|ES`0tg_000IagfB*srAh4al{{tY{IjsNy literal 0 HcmV?d00001 diff --git a/.todos/issues.db-shm b/.todos/issues.db-shm new file mode 100644 index 0000000000000000000000000000000000000000..167349fd8be198e8de77199cf34b6c05377f1939 GIT binary patch literal 32768 zcmeI)Niu^`5Ww+&%!CMnn1>*yAf^y8hnR^g*eJ(w0A=k6mX6{K_IzJf*(z08sDD@W zdtF_x`c?f7(8l1{xRfB8Nf6-GGs zyDBQps$(Fi&A%H|ky7>3Mim}0tzUgfC36Apnw7jD4>7>3Mim} z0tzUgfC36Apnw7jD4>7>3Mim}0tzUgfC36Apnw7jD4>7>3Mim}0tzUgfC36Apnw83 z5QtDmJq7>3Mim}0tzUgfC36Apnw7jD4>7>3Mim}0tzUgfC36Apnw7j YD4>7>3Mim}0tzUgfC36A@DBuj03kd#3;+NC literal 0 HcmV?d00001 diff --git a/.todos/issues.db-wal b/.todos/issues.db-wal new file mode 100644 index 0000000000000000000000000000000000000000..a0287ce218a459598118a90b8514e2e868104755 GIT binary patch literal 337872 zcmeI*f2^KmeaG>qlh(FS3k%)ABAk+x<#QVO_U;5m; z|G3`-{ez@ea=o7KKmFu;eeSq*(yg8T@1N1>oZ9LBbHVN>uRnd!sznz*^vGquckJoO z@}JI>|1JK%|Ng!H10UbG>yF!|$$y+MeQ8JjG*A3~Ivg;dN zwRZE`A#Ici{mS%f)2HpA_k$`7Wanazee;x4s{8+a>xSm30d2JoIfPjKBXK}!=d3OM zkuVUQ(;>Q`*B#(Z}FINb(aEa)G~p>iTP+f97Wk*cV`U5kLR|1Q0*~0R#|0 z009IL_{{|x-WSLNDrpx>D+XyGX>GJgF~|)Ixxn@(wr>2>PcD7-L@uzTGkwW#Zoml$ zAb?0`wvY?F@zQ^IY}0G`QgVS_jsa&OfB*srAbF%a zdIz0yS(|L-cPZoo58b)r$+10$w~`CAbYi(S0tg_000IagfB*srAb>zG1)9wTjFx_Z z98)Of0!DfVZT4ooeSw{~|M<*9Z=Lb|6S=_B&h(|d5I_I{1Q0*~0R#|0009JA zB{0_PcM6S_hXh`ga{=icG%m?Yg6ek)|MtcHE6@DT{d3qCXw{r@Sp*P3009ILKmY** z5I_KddITEY7x3M81kT&izJT-&I;FGALCLiX?+Da_JCC1!<$^`z0`&|8ry_s=0tg_0 z00IagfB*srv`V1iT%g-Kn4Qa|T%g-Kn03}sS8{>h|Lt!*eaBD!?GMQXS~aI!76Akh zKmY**5I_I{1Q0-=9)X5)fh@T|w5jwx0wcYHCTeTO`A*@*f3xi^C(pU%9CCqr27*%& zKmY**5I_I{1Q0*~0R&nl&}=SXwd4XeDqq|eu+lqdtn*5_%J&GK-}2=>w>)>%d&vb_ zHK$w_0R#|0009ILKmY**5I~?FfrfJdEByk_SY65mq<7HS9K$&K0#DxgC+}M}`-!9E z0`&|8ry_s=0tg_000IagfB*srv`V1iT%g-KXq_#+Be31xL93GsCdjW*=pEd4|9kh$ zy=U-`$pu<9r(6~R1Q0*~0R#|0009ILK%gFhhI4_~%>{CZrF{YE9dtIzCV|&-?Lsc_ z?KJ1rrJKLBhFqYYf#6gG5I_I{1Q0*~0R#|00D)EsG@J``dk0fYrSBBlZtr07R)>(} z*C^xy6F>aps~*1Xz+L15t(sFVivR)$Ab;as5GI~ap0 zeW%cMdk3S9K34k$Ubyut=s(l4;$+oylQ1Q0*~0R#|0009IL zK%gdphI0Yw9h46eI9J>k@X|YIbI2)IbAhjH{Ms{T?fueb_62I12#!Sn0R#|0009IL zKmY**5NMM?!?}R;4$9{Ub9rAtdIx>--b7v57r6TC5551=ciuITT%b*3%3Tpa009IL zKmY**5I_I{1Zom!I2VxKLG84)rF{YE9dtIyx_}9C@4`ERUktr(@w)ivTylY$CW2!T zKmY**5I_I{1Q0*~0R-A4&~Pr$?HzQ^n^G>&?HzPRjrWdV_CB-hj`KHvfLx$WW6E6- zKmY**5I_I{1Q0*~0R(CiXfhWFgGzb_P0p(Hjvz?ypwA)rn62EqurF}Q(Z9Un+`*>| zxj;=5!LbM+fB*srAb*!KZ5mVViU0x#Ab=v`DZ21^TtzNW(?oDA0tg_000IagfB*srAb>!d z1RBl-q<7F*8$h`~i2;T8m54hjDyzgFeftn_QV-Y|A0R#|0 z009ILKmY**+9c3$E+D;wHW^cTN06jSY0AQz}< zA~+TS1Q0*~0R#|0009ILK%h+m4d(*VJLpWv(>vxxnkagE5!-1@h~? zgFfVOdIzVUzHa-QZur7KlMB={5S)qt0tg_000IagfB*srAkZp-hI4^#@1R$y+%M4W z9dyZORqYqpbKcJvUN-;g87FdqYdh1gZPo9^Wf4FC0R#|0009ILKmY**5a^-6*f|qU zo;7X8%*peojdr$-t{h#mTrXd_EUa6$dU(_Dn&I_ZhDSE6Uoki`viiCq9aDCFqpQ|# zUOS|XGNE6Yer@`+9rS)sg@No`%%wL0O8(}rb;u#a(jN(fu3I}8Qn_A0E9o7KAp~b_ zrC;ER{;%GA#ole6T%d&O*ca&G7;qE<2q1s}0tg_000IagfIv$H8r~O> z+QE?JwLm!+klw+Zjl5N;ey4ECl)Ha+&zg--vM%P5AbBw?z1Q0*~0R#|0 z009ILKmdUr3N)MxNbjJJE|%XBXz3kH(Re>j@8C`Uea)oDpIb12T%d??j^N%?(mAs(S-qcJprs?r zwGlu70R#|0009ILKmY**dMMCvE+D;w(V9xXfDZDGz^Z7Su6&PR;{_L-bbsgKZ;}i2 za11yK0R#|0009ILKmY**5I~@%0uAQ^Qak8!OqG2B=^c#Ohd9nVf*lXu{zqT?&cU0= z1zI|?TpIxd5I_I{1Q0*~0R#|0poaoY<^pk0N$+5CPKQ$Oppo7|?TpgiRo)RC+|1*CT{#h5DZ z2$JLi(MRWklFJryfk!`a<=n?Uv*4@b0zDi9jzRzd1Q0*~0R#|0009ILXsJN6xqz15 z!5D%ozay~HJE(%QX`Js7^u2iFz5DOEZz;JzOGlP#BY*$`2q1s}0tg_000Ic~P@v&l zpxZm>e0HUG1h(5d*o_5J4074RzCifF&AUFY-ZkSyF0iCCeMt{FI0^v-5I_I{1Q0*~ z0R#|00D)EsjE$W=;pADBBzcT&W^l3Zj{h*Ro0Ow+k#lPjN0d2JoIfPjHBXLkm z@nCeZ{GCE8y@SemZ-S}h0>?Hy|LDK&`)FiepjCs)Wf4FC0R#|0009ILKmY**Y7uC7 zUqE^Xy@|f^jvx%AWW3igTe)^&U*P=k^1tkRZN+@{1!@@x4n+U~1Q0*~0R#|0009IL zXq7r)SXYW7O@1zI(zTowTY5I_I{1Q0*~ z0R#|0pca8<_XUjf4t5_5lyd_<{mWc3$$ubxhw(* zAbJF@GOSFtZp z%Rq1_0tg_000IagfB*srAb>!t1e)F#kmA9tQso_ieA&RL;Nv*^0^hiB=c}Lo)=v+z zFVLz%<+2DMfB*srAbOx!NzV zc= zl6M3-2OUh1pQDfq?40=6FP{0w%MXzYv}#bfECL81fB*srAb&{!#mzt6+94OHWgs{d0R#|0009ILKmY**5I~?+0uAQ^ zQahM~OXYV2!F1mdm~6*+NAQi8ZvF0>2X}v#T%c8h%4HEi009ILKmY**5I_I{1Zojz zI2VxK!C+D-e~%!zfh55|$kA)LcHtetok!l(8N4rQa)DX~f, pub workspace: Option>, pub extension_manager: Option>, + pub hooks: Arc, } /// The main agent that coordinates all components. @@ -116,6 +118,7 @@ impl Agent { deps.safety.clone(), deps.tools.clone(), deps.store.clone(), + deps.hooks.clone(), )); Self { @@ -158,6 +161,10 @@ impl Agent { self.deps.workspace.as_ref() } + fn hooks(&self) -> &Arc { + &self.deps.hooks + } + /// Run the agent main loop. pub async fn run(self) -> Result<(), Error> { // Start channels @@ -425,10 +432,32 @@ impl Agent { match self.handle_message(&message).await { Ok(Some(response)) if !response.is_empty() => { - let _ = self - .channels - .respond(&message, OutgoingResponse::text(response)) - .await; + // Hook: BeforeOutbound — allow hooks to modify or suppress outbound + let event = crate::hooks::HookEvent::Outbound { + user_id: message.user_id.clone(), + channel: message.channel.clone(), + content: response.clone(), + thread_id: message.thread_id.clone(), + }; + match self.hooks().run(&event).await { + Err(err) => { + tracing::warn!("BeforeOutbound hook blocked response: {}", err); + } + Ok(crate::hooks::HookOutcome::Continue { + modified: Some(new_content), + }) => { + let _ = self + .channels + .respond(&message, OutgoingResponse::text(new_content)) + .await; + } + _ => { + let _ = self + .channels + .respond(&message, OutgoingResponse::text(response)) + .await; + } + } } Ok(Some(_)) => { // Empty response, nothing to send (e.g. approval handled via send_status) @@ -474,7 +503,33 @@ impl Agent { async fn handle_message(&self, message: &IncomingMessage) -> Result, Error> { // Parse submission type first - let submission = SubmissionParser::parse(&message.content); + let mut submission = SubmissionParser::parse(&message.content); + + // Hook: BeforeInbound — allow hooks to modify or reject user input + if let Submission::UserInput { ref content } = submission { + let event = crate::hooks::HookEvent::Inbound { + user_id: message.user_id.clone(), + channel: message.channel.clone(), + content: content.clone(), + thread_id: message.thread_id.clone(), + }; + match self.hooks().run(&event).await { + Err(crate::hooks::HookError::Rejected { reason }) => { + return Ok(Some(format!("[Message rejected: {}]", reason))); + } + Err(err) => { + return Ok(Some(format!("[Message blocked by hook policy: {}]", err))); + } + Ok(crate::hooks::HookOutcome::Continue { + modified: Some(new_content), + }) => { + submission = Submission::UserInput { + content: new_content, + }; + } + _ => {} // Continue, fail-open errors already logged in registry + } + } // Hydrate thread from DB if it's a historical thread not in memory if let Some(ref external_thread_id) = message.thread_id { @@ -883,6 +938,27 @@ impl Agent { // Complete, fail, or request approval match result { Ok(AgenticLoopResult::Response(response)) => { + // Hook: TransformResponse — allow hooks to modify or reject the final response + let response = { + let event = crate::hooks::HookEvent::ResponseTransform { + user_id: message.user_id.clone(), + thread_id: thread_id.to_string(), + response: response.clone(), + }; + match self.hooks().run(&event).await { + Err(crate::hooks::HookError::Rejected { reason }) => { + format!("[Response filtered: {}]", reason) + } + Err(err) => { + format!("[Response blocked by hook policy: {}]", err) + } + Ok(crate::hooks::HookOutcome::Continue { + modified: Some(new_response), + }) => new_response, + _ => response, // fail-open: use original + } + }; + thread.complete_turn(&response); self.persist_response_chain(thread); let _ = self @@ -1160,8 +1236,8 @@ impl Agent { } } - // Execute each tool (with approval checking) - for tc in tool_calls { + // Execute each tool (with approval checking and hook interception) + for mut tc in tool_calls { // Check if tool requires approval if let Some(tool) = self.tools().get(&tc.name).await && tool.requires_approval() @@ -1216,6 +1292,47 @@ impl Agent { } } + // Hook: BeforeToolCall — allow hooks to modify or reject tool calls + { + let event = crate::hooks::HookEvent::ToolCall { + tool_name: tc.name.clone(), + parameters: tc.arguments.clone(), + user_id: message.user_id.clone(), + context: "chat".to_string(), + }; + match self.hooks().run(&event).await { + Err(crate::hooks::HookError::Rejected { reason }) => { + context_messages.push(ChatMessage::tool_result( + &tc.id, + &tc.name, + format!("Tool call rejected by hook: {}", reason), + )); + continue; + } + Err(err) => { + context_messages.push(ChatMessage::tool_result( + &tc.id, + &tc.name, + format!("Tool call blocked by hook policy: {}", err), + )); + continue; + } + Ok(crate::hooks::HookOutcome::Continue { + modified: Some(new_params), + }) => match serde_json::from_str(&new_params) { + Ok(parsed) => tc.arguments = parsed, + Err(e) => { + tracing::warn!( + tool = %tc.name, + "Hook returned non-JSON modification for ToolCall, ignoring: {}", + e + ); + } + }, + _ => {} // Continue, fail-open errors already logged + } + } + let _ = self .channels .send_status( diff --git a/src/agent/scheduler.rs b/src/agent/scheduler.rs index a665c2af..23b9ea7c 100644 --- a/src/agent/scheduler.rs +++ b/src/agent/scheduler.rs @@ -14,6 +14,7 @@ use crate::config::AgentConfig; use crate::context::{ContextManager, JobContext, JobState}; use crate::db::Database; use crate::error::{Error, JobError}; +use crate::hooks::HookRegistry; use crate::llm::LlmProvider; use crate::safety::SafetyLayer; use crate::tools::ToolRegistry; @@ -49,6 +50,7 @@ pub struct Scheduler { safety: Arc, tools: Arc, store: Option>, + hooks: Arc, /// Running jobs (main LLM-driven jobs). jobs: Arc>>, /// Running sub-tasks (tool executions, background tasks). @@ -64,6 +66,7 @@ impl Scheduler { safety: Arc, tools: Arc, store: Option>, + hooks: Arc, ) -> Self { Self { config, @@ -72,6 +75,7 @@ impl Scheduler { safety, tools, store, + hooks, jobs: Arc::new(RwLock::new(HashMap::new())), subtasks: Arc::new(RwLock::new(HashMap::new())), } @@ -118,6 +122,7 @@ impl Scheduler { safety: self.safety.clone(), tools: self.tools.clone(), store: self.store.clone(), + hooks: self.hooks.clone(), timeout: self.config.job_timeout, use_planning: self.config.use_planning, }; diff --git a/src/agent/session_manager.rs b/src/agent/session_manager.rs index db0be886..244348cd 100644 --- a/src/agent/session_manager.rs +++ b/src/agent/session_manager.rs @@ -11,6 +11,7 @@ use uuid::Uuid; use crate::agent::session::Session; use crate::agent::undo::UndoManager; +use crate::hooks::HookRegistry; /// Key for mapping external thread IDs to internal ones. #[derive(Clone, Hash, Eq, PartialEq)] @@ -25,6 +26,7 @@ pub struct SessionManager { sessions: RwLock>>>, thread_map: RwLock>, undo_managers: RwLock>>>, + hooks: Option>, } impl SessionManager { @@ -34,9 +36,16 @@ impl SessionManager { sessions: RwLock::new(HashMap::new()), thread_map: RwLock::new(HashMap::new()), undo_managers: RwLock::new(HashMap::new()), + hooks: None, } } + /// Attach a hook registry for session lifecycle events. + pub fn with_hooks(mut self, hooks: Arc) -> Self { + self.hooks = Some(hooks); + self + } + /// Get or create a session for a user. pub async fn get_or_create_session(&self, user_id: &str) -> Arc> { // Fast path: check if session exists @@ -54,8 +63,28 @@ impl SessionManager { return Arc::clone(session); } - let session = Arc::new(Mutex::new(Session::new(user_id))); + let new_session = Session::new(user_id); + let session_id = new_session.id.to_string(); + let session = Arc::new(Mutex::new(new_session)); sessions.insert(user_id.to_string(), Arc::clone(&session)); + + // Fire OnSessionStart hook (fire-and-forget) + if let Some(ref hooks) = self.hooks { + let hooks = hooks.clone(); + let uid = user_id.to_string(); + let sid = session_id; + tokio::spawn(async move { + use crate::hooks::HookEvent; + let event = HookEvent::SessionStart { + user_id: uid, + session_id: sid, + }; + if let Err(e) = hooks.run(&event).await { + tracing::warn!("OnSessionStart hook error: {}", e); + } + }); + } + session } @@ -173,8 +202,8 @@ impl SessionManager { pub async fn prune_stale_sessions(&self, max_idle: std::time::Duration) -> usize { let cutoff = chrono::Utc::now() - chrono::TimeDelta::seconds(max_idle.as_secs() as i64); - // Find stale session user_ids - let stale_users: Vec = { + // Find stale sessions (user_id + session_id) + let stale_sessions: Vec<(String, String)> = { let sessions = self.sessions.read().await; sessions .iter() @@ -182,7 +211,7 @@ impl SessionManager { // Try to lock; skip if contended (someone is actively using it) let sess = session.try_lock().ok()?; if sess.last_active_at < cutoff { - Some(user_id.clone()) + Some((user_id.clone(), sess.id.to_string())) } else { None } @@ -190,6 +219,11 @@ impl SessionManager { .collect() }; + let stale_users: Vec = stale_sessions + .iter() + .map(|(user_id, _)| user_id.clone()) + .collect(); + if stale_users.is_empty() { return 0; } @@ -207,6 +241,25 @@ impl SessionManager { } } + // Fire OnSessionEnd hooks for stale sessions (fire-and-forget) + if let Some(ref hooks) = self.hooks { + for (user_id, session_id) in &stale_sessions { + let hooks = hooks.clone(); + let uid = user_id.clone(); + let sid = session_id.clone(); + tokio::spawn(async move { + use crate::hooks::HookEvent; + let event = HookEvent::SessionEnd { + user_id: uid, + session_id: sid, + }; + if let Err(e) = hooks.run(&event).await { + tracing::warn!("OnSessionEnd hook error: {}", e); + } + }); + } + } + // Remove sessions let count = { let mut sessions = self.sessions.write().await; diff --git a/src/agent/worker.rs b/src/agent/worker.rs index 565a3d22..3ec88586 100644 --- a/src/agent/worker.rs +++ b/src/agent/worker.rs @@ -12,6 +12,7 @@ use crate::agent::task::TaskOutput; use crate::context::{ContextManager, JobState}; use crate::db::Database; use crate::error::Error; +use crate::hooks::HookRegistry; use crate::llm::{ ActionPlan, ChatMessage, LlmProvider, Reasoning, ReasoningContext, RespondResult, ToolSelection, }; @@ -29,6 +30,7 @@ pub struct WorkerDeps { pub safety: Arc, pub tools: Arc, pub store: Option>, + pub hooks: Arc, pub timeout: Duration, pub use_planning: bool, } @@ -352,23 +354,11 @@ Report when the job is complete or if you encounter issues you cannot resolve."# .map(|selection| { let tool_name = selection.tool_name.clone(); let params = selection.parameters.clone(); - let tools = self.tools().clone(); - let context_manager = self.context_manager().clone(); - let safety = self.safety().clone(); + let deps = self.deps.clone(); let job_id = self.job_id; - let store = self.deps.store.clone(); async move { - let result = Self::execute_tool_inner( - tools, - context_manager, - safety, - store, - job_id, - &tool_name, - ¶ms, - ) - .await; + let result = Self::execute_tool_inner(&deps, job_id, &tool_name, ¶ms).await; ToolExecResult { result } } }) @@ -379,20 +369,18 @@ Report when the job is complete or if you encounter issues you cannot resolve."# /// Inner tool execution logic that can be called from both single and parallel paths. async fn execute_tool_inner( - tools: Arc, - context_manager: Arc, - safety: Arc, - store: Option>, + deps: &WorkerDeps, job_id: Uuid, tool_name: &str, params: &serde_json::Value, ) -> Result { - let tool = tools - .get(tool_name) - .await - .ok_or_else(|| crate::error::ToolError::NotFound { - name: tool_name.to_string(), - })?; + let tool = + deps.tools + .get(tool_name) + .await + .ok_or_else(|| crate::error::ToolError::NotFound { + name: tool_name.to_string(), + })?; // Tools requiring approval are blocked in autonomous jobs if tool.requires_approval() { @@ -402,8 +390,46 @@ Report when the job is complete or if you encounter issues you cannot resolve."# .into()); } - // Get job context for the tool - let job_ctx = context_manager.get_context(job_id).await?; + // Fetch job context early so we have the real user_id for hooks + let job_ctx = deps.context_manager.get_context(job_id).await?; + + // Run BeforeToolCall hook + let params = { + use crate::hooks::{HookError, HookEvent, HookOutcome}; + let event = HookEvent::ToolCall { + tool_name: tool_name.to_string(), + parameters: params.clone(), + user_id: job_ctx.user_id.clone(), + context: format!("job:{}", job_id), + }; + match deps.hooks.run(&event).await { + Err(HookError::Rejected { reason }) => { + return Err(crate::error::ToolError::ExecutionFailed { + name: tool_name.to_string(), + reason: format!("Blocked by hook: {}", reason), + } + .into()); + } + Err(err) => { + return Err(crate::error::ToolError::ExecutionFailed { + name: tool_name.to_string(), + reason: format!("Blocked by hook failure mode: {}", err), + } + .into()); + } + Ok(HookOutcome::Continue { + modified: Some(new_params), + }) => serde_json::from_str(&new_params).unwrap_or_else(|e| { + tracing::warn!( + tool = %tool_name, + "Hook returned non-JSON modification for ToolCall, ignoring: {}", + e + ); + params.clone() + }), + _ => params.clone(), + } + }; if job_ctx.state == JobState::Cancelled { return Err(crate::error::ToolError::ExecutionFailed { name: tool_name.to_string(), @@ -413,7 +439,7 @@ Report when the job is complete or if you encounter issues you cannot resolve."# } // Validate tool parameters - let validation = safety.validator().validate_tool_params(params); + let validation = deps.safety.validator().validate_tool_params(¶ms); if !validation.is_valid { let details = validation .errors @@ -478,8 +504,8 @@ Report when the job is complete or if you encounter issues you cannot resolve."# Ok(Ok(output)) => { let output_str = serde_json::to_string_pretty(&output.result) .ok() - .map(|s| safety.sanitize_tool_output(tool_name, &s).content); - context_manager + .map(|s| deps.safety.sanitize_tool_output(tool_name, &s).content); + deps.context_manager .update_memory(job_id, |mem| { let rec = mem.create_action(tool_name, params.clone()).succeed( output_str.clone(), @@ -492,7 +518,8 @@ Report when the job is complete or if you encounter issues you cannot resolve."# .await .ok() } - Ok(Err(e)) => context_manager + Ok(Err(e)) => deps + .context_manager .update_memory(job_id, |mem| { let rec = mem .create_action(tool_name, params.clone()) @@ -502,7 +529,8 @@ Report when the job is complete or if you encounter issues you cannot resolve."# }) .await .ok(), - Err(_) => context_manager + Err(_) => deps + .context_manager .update_memory(job_id, |mem| { let rec = mem .create_action(tool_name, params.clone()) @@ -515,7 +543,7 @@ Report when the job is complete or if you encounter issues you cannot resolve."# }; // Persist action to database (fire-and-forget) - if let (Some(action), Some(store)) = (action, store) { + if let (Some(action), Some(store)) = (action, deps.store.clone()) { tokio::spawn(async move { if let Err(e) = store.save_action(job_id, &action).await { tracing::warn!("Failed to persist action for job {}: {}", job_id, e); @@ -701,16 +729,7 @@ Report when the job is complete or if you encounter issues you cannot resolve."# tool_name: &str, params: &serde_json::Value, ) -> Result { - Self::execute_tool_inner( - self.tools().clone(), - self.context_manager().clone(), - self.safety().clone(), - self.deps.store.clone(), - self.job_id, - tool_name, - params, - ) - .await + Self::execute_tool_inner(&self.deps, self.job_id, tool_name, params).await } async fn mark_completed(&self) -> Result<(), Error> { diff --git a/src/error.rs b/src/error.rs index acb3598a..5aa7d461 100644 --- a/src/error.rs +++ b/src/error.rs @@ -40,6 +40,9 @@ pub enum Error { #[error("Workspace error: {0}")] Workspace(#[from] WorkspaceError), + #[error("Hook error: {0}")] + Hook(#[from] crate::hooks::HookError), + #[error("Orchestrator error: {0}")] Orchestrator(#[from] OrchestratorError), diff --git a/src/hooks/hook.rs b/src/hooks/hook.rs new file mode 100644 index 00000000..7df174e1 --- /dev/null +++ b/src/hooks/hook.rs @@ -0,0 +1,199 @@ +//! Core hook types and traits. + +use std::time::Duration; + +use async_trait::async_trait; + +/// Points in the agent lifecycle where hooks can be attached. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub enum HookPoint { + /// Before processing an inbound user message. + BeforeInbound, + /// Before executing a tool call. + BeforeToolCall, + /// Before sending an outbound response. + BeforeOutbound, + /// When a new session starts. + OnSessionStart, + /// When a session ends (pruned or expired). + OnSessionEnd, + /// Transform the final response before completing a turn. + TransformResponse, +} + +/// Contextual data carried with each hook invocation. +#[derive(Debug, Clone)] +pub enum HookEvent { + /// An inbound user message about to be processed. + Inbound { + user_id: String, + channel: String, + content: String, + thread_id: Option, + }, + /// A tool call about to be executed. + ToolCall { + tool_name: String, + parameters: serde_json::Value, + user_id: String, + /// "chat" for interactive, or a job ID string for autonomous jobs. + context: String, + }, + /// An outbound response about to be sent. + Outbound { + user_id: String, + channel: String, + content: String, + thread_id: Option, + }, + /// A new session was created. + SessionStart { user_id: String, session_id: String }, + /// A session was ended (pruned). + SessionEnd { user_id: String, session_id: String }, + /// The final response is being transformed before completing a turn. + ResponseTransform { + user_id: String, + thread_id: String, + response: String, + }, +} + +impl HookEvent { + /// Returns the [`HookPoint`] this event corresponds to. + pub fn hook_point(&self) -> HookPoint { + match self { + HookEvent::Inbound { .. } => HookPoint::BeforeInbound, + HookEvent::ToolCall { .. } => HookPoint::BeforeToolCall, + HookEvent::Outbound { .. } => HookPoint::BeforeOutbound, + HookEvent::SessionStart { .. } => HookPoint::OnSessionStart, + HookEvent::SessionEnd { .. } => HookPoint::OnSessionEnd, + HookEvent::ResponseTransform { .. } => HookPoint::TransformResponse, + } + } + + /// Apply a modification string to the event's primary content field. + pub fn apply_modification(&mut self, modified: &str) { + match self { + HookEvent::Inbound { content, .. } | HookEvent::Outbound { content, .. } => { + *content = modified.to_string(); + } + HookEvent::ToolCall { parameters, .. } => match serde_json::from_str(modified) { + Ok(parsed) => *parameters = parsed, + Err(e) => { + tracing::warn!( + "Hook returned non-JSON modification for ToolCall, ignoring: {}", + e + ); + } + }, + HookEvent::ResponseTransform { response, .. } => { + *response = modified.to_string(); + } + HookEvent::SessionStart { .. } | HookEvent::SessionEnd { .. } => { + // Session events don't have modifiable content + } + } + } +} + +/// The result of executing a hook. +#[derive(Debug, Clone)] +pub enum HookOutcome { + /// Continue processing, optionally with modified content. + Continue { + /// If `Some`, replace the event's primary content with this value. + modified: Option, + }, + /// Reject the event entirely. + Reject { + /// Human-readable reason for the rejection. + reason: String, + }, +} + +impl HookOutcome { + /// Shorthand for `Continue { modified: None }`. + pub fn ok() -> Self { + HookOutcome::Continue { modified: None } + } + + /// Shorthand for `Continue { modified: Some(value) }`. + pub fn modify(value: String) -> Self { + HookOutcome::Continue { + modified: Some(value), + } + } + + /// Shorthand for `Reject { reason }`. + pub fn reject(reason: impl Into) -> Self { + HookOutcome::Reject { + reason: reason.into(), + } + } +} + +/// How to handle hook execution failures. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum HookFailureMode { + /// On error/timeout, continue processing as if the hook returned `ok()`. + FailOpen, + /// On error/timeout, reject the event. + FailClosed, +} + +/// Hook execution errors. +#[derive(Debug, thiserror::Error)] +pub enum HookError { + #[error("Hook execution failed: {reason}")] + ExecutionFailed { reason: String }, + + #[error("Hook timed out after {timeout:?}")] + Timeout { timeout: Duration }, + + #[error("Hook rejected: {reason}")] + Rejected { reason: String }, +} + +/// Context passed to hooks alongside the event. +pub struct HookContext { + /// Arbitrary metadata hooks can use. + pub metadata: serde_json::Value, +} + +impl Default for HookContext { + fn default() -> Self { + Self { + metadata: serde_json::Value::Null, + } + } +} + +/// Trait for implementing lifecycle hooks. +/// +/// Hooks intercept and can modify agent operations at well-defined points. +#[async_trait] +pub trait Hook: Send + Sync { + /// A unique name for this hook. + fn name(&self) -> &str; + + /// The lifecycle points this hook should be called at. + fn hook_points(&self) -> &[HookPoint]; + + /// How to handle failures in this hook. + /// + /// Default: `FailOpen` (continue on error). + fn failure_mode(&self) -> HookFailureMode { + HookFailureMode::FailOpen + } + + /// Maximum time this hook is allowed to run. + /// + /// Default: 5 seconds. + fn timeout(&self) -> Duration { + Duration::from_secs(5) + } + + /// Execute the hook. + async fn execute(&self, event: &HookEvent, ctx: &HookContext) + -> Result; +} diff --git a/src/hooks/mod.rs b/src/hooks/mod.rs new file mode 100644 index 00000000..9ea6a5ac --- /dev/null +++ b/src/hooks/mod.rs @@ -0,0 +1,19 @@ +//! Lifecycle hooks for intercepting and transforming agent operations. +//! +//! The hook system provides 6 well-defined interception points: +//! +//! - **BeforeInbound** — Before processing an inbound user message +//! - **BeforeToolCall** — Before executing a tool call +//! - **BeforeOutbound** — Before sending an outbound response +//! - **OnSessionStart** — When a new session starts +//! - **OnSessionEnd** — When a session ends +//! - **TransformResponse** — Transform the final response before completing a turn +//! +//! Hooks are executed in priority order (lower number = higher priority). +//! Each hook can pass through, modify content, or reject the event. + +pub mod hook; +pub mod registry; + +pub use hook::{Hook, HookContext, HookError, HookEvent, HookFailureMode, HookOutcome, HookPoint}; +pub use registry::HookRegistry; diff --git a/src/hooks/registry.rs b/src/hooks/registry.rs new file mode 100644 index 00000000..6148d954 --- /dev/null +++ b/src/hooks/registry.rs @@ -0,0 +1,555 @@ +//! Hook registry for managing and executing lifecycle hooks. + +use std::sync::Arc; + +use tokio::sync::RwLock; + +use crate::hooks::hook::{Hook, HookContext, HookError, HookEvent, HookFailureMode, HookOutcome}; + +/// A registered hook with its priority. +struct HookEntry { + hook: Arc, + priority: u32, +} + +/// Registry that manages hooks and executes them at lifecycle points. +/// +/// Hooks are executed in priority order (lower number = higher priority). +/// A `Reject` outcome stops the chain immediately. +/// A `Modify` outcome chains through subsequent hooks. +pub struct HookRegistry { + hooks: RwLock>, +} + +impl HookRegistry { + /// Create an empty registry. + pub fn new() -> Self { + Self { + hooks: RwLock::new(Vec::new()), + } + } + + /// Register a hook with default priority (100). + pub async fn register(&self, hook: Arc) { + self.register_with_priority(hook, 100).await; + } + + /// Register a hook with a specific priority. + /// + /// Lower priority number = runs first. + pub async fn register_with_priority(&self, hook: Arc, priority: u32) { + let mut hooks = self.hooks.write().await; + hooks.push(HookEntry { hook, priority }); + hooks.sort_by_key(|e| e.priority); + } + + /// Unregister a hook by name. Returns `true` if it was found and removed. + pub async fn unregister(&self, name: &str) -> bool { + let mut hooks = self.hooks.write().await; + let before = hooks.len(); + hooks.retain(|e| e.hook.name() != name); + hooks.len() < before + } + + /// List all registered hook names (in priority order). + pub async fn list(&self) -> Vec { + let hooks = self.hooks.read().await; + hooks.iter().map(|e| e.hook.name().to_string()).collect() + } + + /// Run all hooks matching the event's hook point. + /// + /// - Hooks run in priority order (lowest first). + /// - `Reject` stops the chain immediately. + /// - `Modify` chains the modification through subsequent hooks. + /// - Timeout/error handling respects each hook's `failure_mode`. + pub async fn run(&self, event: &HookEvent) -> Result { + let point = event.hook_point(); + let ctx = HookContext::default(); + + // Clone matching hooks and drop the read guard before executing. + // Each hook can run up to its timeout, so holding the guard would + // block concurrent register/unregister/run calls. + let matching: Vec> = { + let hooks = self.hooks.read().await; + hooks + .iter() + .filter(|e| e.hook.hook_points().contains(&point)) + .map(|e| e.hook.clone()) + .collect() + }; + + if matching.is_empty() { + return Ok(HookOutcome::ok()); + } + + let mut current_event = event.clone(); + + for hook in &matching { + let timeout = hook.timeout(); + + let result = tokio::time::timeout(timeout, hook.execute(¤t_event, &ctx)).await; + + match result { + Ok(Ok(HookOutcome::Reject { reason })) => { + tracing::debug!(hook = hook.name(), "Hook rejected: {}", reason); + return Err(HookError::Rejected { reason }); + } + Ok(Ok(HookOutcome::Continue { + modified: Some(value), + })) => { + tracing::debug!(hook = hook.name(), "Hook modified content"); + current_event.apply_modification(&value); + } + Ok(Ok(HookOutcome::Continue { modified: None })) => { + // No-op, continue chain + } + Ok(Err(err)) => match hook.failure_mode() { + HookFailureMode::FailOpen => { + tracing::warn!(hook = hook.name(), "Hook failed (fail-open): {}", err); + } + HookFailureMode::FailClosed => { + tracing::warn!(hook = hook.name(), "Hook failed (fail-closed): {}", err); + return Err(HookError::ExecutionFailed { + reason: format!("Hook '{}' failed: {}", hook.name(), err), + }); + } + }, + Err(_elapsed) => match hook.failure_mode() { + HookFailureMode::FailOpen => { + tracing::warn!( + hook = hook.name(), + "Hook timed out (fail-open) after {:?}", + timeout + ); + } + HookFailureMode::FailClosed => { + tracing::warn!( + hook = hook.name(), + "Hook timed out (fail-closed) after {:?}", + timeout + ); + return Err(HookError::Timeout { timeout }); + } + }, + } + } + + // Determine final outcome by comparing with original event + let modified = extract_content(¤t_event); + let original = extract_content(event); + + if modified != original { + Ok(HookOutcome::modify(modified)) + } else { + Ok(HookOutcome::ok()) + } + } +} + +impl Default for HookRegistry { + fn default() -> Self { + Self::new() + } +} + +/// Extract the primary content string from a hook event. +fn extract_content(event: &HookEvent) -> String { + match event { + HookEvent::Inbound { content, .. } | HookEvent::Outbound { content, .. } => content.clone(), + HookEvent::ToolCall { parameters, .. } => { + serde_json::to_string(parameters).unwrap_or_default() + } + HookEvent::ResponseTransform { response, .. } => response.clone(), + HookEvent::SessionStart { session_id, .. } | HookEvent::SessionEnd { session_id, .. } => { + session_id.clone() + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::hooks::hook::{HookFailureMode, HookPoint}; + use async_trait::async_trait; + use std::time::Duration; + + /// A test hook that always returns ok. + struct PassthroughHook { + name: String, + points: Vec, + } + + #[async_trait] + impl Hook for PassthroughHook { + fn name(&self) -> &str { + &self.name + } + fn hook_points(&self) -> &[HookPoint] { + &self.points + } + async fn execute( + &self, + _event: &HookEvent, + _ctx: &HookContext, + ) -> Result { + Ok(HookOutcome::ok()) + } + } + + /// A hook that modifies content by appending a suffix. + struct ModifyHook { + name: String, + suffix: String, + points: Vec, + } + + #[async_trait] + impl Hook for ModifyHook { + fn name(&self) -> &str { + &self.name + } + fn hook_points(&self) -> &[HookPoint] { + &self.points + } + async fn execute( + &self, + event: &HookEvent, + _ctx: &HookContext, + ) -> Result { + let content = extract_content(event); + Ok(HookOutcome::modify(format!("{}{}", content, self.suffix))) + } + } + + /// A hook that always rejects. + struct RejectHook { + name: String, + reason: String, + points: Vec, + } + + #[async_trait] + impl Hook for RejectHook { + fn name(&self) -> &str { + &self.name + } + fn hook_points(&self) -> &[HookPoint] { + &self.points + } + async fn execute( + &self, + _event: &HookEvent, + _ctx: &HookContext, + ) -> Result { + Ok(HookOutcome::reject(&self.reason)) + } + } + + /// A hook that always errors. + struct ErrorHook { + name: String, + points: Vec, + failure_mode: HookFailureMode, + } + + #[async_trait] + impl Hook for ErrorHook { + fn name(&self) -> &str { + &self.name + } + fn hook_points(&self) -> &[HookPoint] { + &self.points + } + fn failure_mode(&self) -> HookFailureMode { + self.failure_mode + } + async fn execute( + &self, + _event: &HookEvent, + _ctx: &HookContext, + ) -> Result { + Err(HookError::ExecutionFailed { + reason: "test error".into(), + }) + } + } + + /// A hook that sleeps longer than its timeout. + struct SlowHook { + name: String, + points: Vec, + failure_mode: HookFailureMode, + } + + #[async_trait] + impl Hook for SlowHook { + fn name(&self) -> &str { + &self.name + } + fn hook_points(&self) -> &[HookPoint] { + &self.points + } + fn failure_mode(&self) -> HookFailureMode { + self.failure_mode + } + fn timeout(&self) -> Duration { + Duration::from_millis(50) + } + async fn execute( + &self, + _event: &HookEvent, + _ctx: &HookContext, + ) -> Result { + tokio::time::sleep(Duration::from_millis(200)).await; + Ok(HookOutcome::ok()) + } + } + + fn test_event() -> HookEvent { + HookEvent::Inbound { + user_id: "user-1".into(), + channel: "test".into(), + content: "hello".into(), + thread_id: None, + } + } + + #[tokio::test] + async fn test_empty_registry_returns_ok() { + let registry = HookRegistry::new(); + let result = registry.run(&test_event()).await; + assert!(result.is_ok()); + assert!(matches!( + result.unwrap(), + HookOutcome::Continue { modified: None } + )); + } + + #[tokio::test] + async fn test_register_and_list() { + let registry = HookRegistry::new(); + registry + .register(Arc::new(PassthroughHook { + name: "hook-a".into(), + points: vec![HookPoint::BeforeInbound], + })) + .await; + registry + .register(Arc::new(PassthroughHook { + name: "hook-b".into(), + points: vec![HookPoint::BeforeInbound], + })) + .await; + + let names = registry.list().await; + assert_eq!(names, vec!["hook-a", "hook-b"]); + } + + #[tokio::test] + async fn test_priority_ordering() { + let registry = HookRegistry::new(); + + // Register in reverse priority order + registry + .register_with_priority( + Arc::new(ModifyHook { + name: "low-prio".into(), + suffix: "-LOW".into(), + points: vec![HookPoint::BeforeInbound], + }), + 200, + ) + .await; + registry + .register_with_priority( + Arc::new(ModifyHook { + name: "high-prio".into(), + suffix: "-HIGH".into(), + points: vec![HookPoint::BeforeInbound], + }), + 10, + ) + .await; + + // Should run in priority order: high-prio first, then low-prio + let names = registry.list().await; + assert_eq!(names[0], "high-prio"); + assert_eq!(names[1], "low-prio"); + + let result = registry.run(&test_event()).await.unwrap(); + match result { + HookOutcome::Continue { modified: Some(m) } => { + // "hello" -> "hello-HIGH" -> "hello-HIGH-LOW" + assert_eq!(m, "hello-HIGH-LOW"); + } + other => panic!("Expected modification chain, got: {:?}", other), + } + } + + #[tokio::test] + async fn test_reject_stops_chain() { + let registry = HookRegistry::new(); + + registry + .register_with_priority( + Arc::new(RejectHook { + name: "blocker".into(), + reason: "blocked".into(), + points: vec![HookPoint::BeforeInbound], + }), + 10, + ) + .await; + registry + .register_with_priority( + Arc::new(ModifyHook { + name: "modifier".into(), + suffix: "-MODIFIED".into(), + points: vec![HookPoint::BeforeInbound], + }), + 20, + ) + .await; + + let result = registry.run(&test_event()).await; + assert!(result.is_err()); + match result.unwrap_err() { + HookError::Rejected { reason } => assert_eq!(reason, "blocked"), + other => panic!("Expected Rejected, got: {:?}", other), + } + } + + #[tokio::test] + async fn test_modification_chaining() { + let registry = HookRegistry::new(); + + registry + .register_with_priority( + Arc::new(ModifyHook { + name: "first".into(), + suffix: "-A".into(), + points: vec![HookPoint::BeforeInbound], + }), + 10, + ) + .await; + registry + .register_with_priority( + Arc::new(ModifyHook { + name: "second".into(), + suffix: "-B".into(), + points: vec![HookPoint::BeforeInbound], + }), + 20, + ) + .await; + + let result = registry.run(&test_event()).await.unwrap(); + match result { + HookOutcome::Continue { modified: Some(m) } => { + assert_eq!(m, "hello-A-B"); + } + other => panic!("Expected chained modification, got: {:?}", other), + } + } + + #[tokio::test] + async fn test_fail_open_on_error() { + let registry = HookRegistry::new(); + registry + .register(Arc::new(ErrorHook { + name: "err-open".into(), + points: vec![HookPoint::BeforeInbound], + failure_mode: HookFailureMode::FailOpen, + })) + .await; + + let result = registry.run(&test_event()).await; + assert!(result.is_ok()); + } + + #[tokio::test] + async fn test_fail_closed_on_error() { + let registry = HookRegistry::new(); + registry + .register(Arc::new(ErrorHook { + name: "err-closed".into(), + points: vec![HookPoint::BeforeInbound], + failure_mode: HookFailureMode::FailClosed, + })) + .await; + + let result = registry.run(&test_event()).await; + assert!(result.is_err()); + assert!(matches!( + result.unwrap_err(), + HookError::ExecutionFailed { .. } + )); + } + + #[tokio::test] + async fn test_fail_open_on_timeout() { + let registry = HookRegistry::new(); + registry + .register(Arc::new(SlowHook { + name: "slow-open".into(), + points: vec![HookPoint::BeforeInbound], + failure_mode: HookFailureMode::FailOpen, + })) + .await; + + let result = registry.run(&test_event()).await; + assert!(result.is_ok()); + } + + #[tokio::test] + async fn test_fail_closed_on_timeout() { + let registry = HookRegistry::new(); + registry + .register(Arc::new(SlowHook { + name: "slow-closed".into(), + points: vec![HookPoint::BeforeInbound], + failure_mode: HookFailureMode::FailClosed, + })) + .await; + + let result = registry.run(&test_event()).await; + assert!(result.is_err()); + assert!(matches!(result.unwrap_err(), HookError::Timeout { .. })); + } + + #[tokio::test] + async fn test_unregister() { + let registry = HookRegistry::new(); + registry + .register(Arc::new(PassthroughHook { + name: "removable".into(), + points: vec![HookPoint::BeforeInbound], + })) + .await; + + assert_eq!(registry.list().await.len(), 1); + assert!(registry.unregister("removable").await); + assert_eq!(registry.list().await.len(), 0); + + // Unregistering non-existent returns false + assert!(!registry.unregister("nonexistent").await); + } + + #[tokio::test] + async fn test_hooks_only_match_their_points() { + let registry = HookRegistry::new(); + registry + .register(Arc::new(RejectHook { + name: "outbound-only".into(), + reason: "blocked".into(), + points: vec![HookPoint::BeforeOutbound], + })) + .await; + + // Inbound event should not be affected by outbound-only hook + let result = registry.run(&test_event()).await; + assert!(result.is_ok()); + } +} diff --git a/src/lib.rs b/src/lib.rs index ed72e3df..2fca3c3f 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -50,6 +50,7 @@ pub mod estimation; pub mod evaluation; pub mod extensions; pub mod history; +pub mod hooks; pub mod llm; pub mod orchestrator; pub mod pairing; diff --git a/src/main.rs b/src/main.rs index fac21c40..1b398def 100644 --- a/src/main.rs +++ b/src/main.rs @@ -22,6 +22,7 @@ use ironclaw::{ config::Config, context::ContextManager, extensions::ExtensionManager, + hooks::HookRegistry, llm::{ FailoverProvider, LlmProvider, SessionConfig, create_cheap_llm_provider, create_llm_provider, create_llm_provider_with_config, create_session_manager, @@ -1132,8 +1133,11 @@ async fn main() -> anyhow::Result<()> { // Create context manager (shared between job tools and agent) let context_manager = Arc::new(ContextManager::new(config.agent.max_parallel_jobs)); + // Create hook registry + let hooks = Arc::new(HookRegistry::new()); + // Create session manager (shared between agent and web gateway) - let session_manager = Arc::new(SessionManager::new()); + let session_manager = Arc::new(SessionManager::new().with_hooks(hooks.clone())); // Register job tools (sandbox deps auto-injected when container_job_manager is available) tools.register_job_tools( @@ -1199,6 +1203,7 @@ async fn main() -> anyhow::Result<()> { tools, workspace, extension_manager, + hooks, }; let agent = Agent::new( config.agent.clone(), diff --git a/src/setup/wizard.rs b/src/setup/wizard.rs index 31c4cc4e..1cc5dd89 100644 --- a/src/setup/wizard.rs +++ b/src/setup/wizard.rs @@ -2090,8 +2090,6 @@ mod tests { #[tokio::test] async fn test_install_missing_bundled_channels_installs_telegram() { - use crate::channels::wasm::available_channel_names; - // WASM artifacts only exist in dev builds (not CI). Skip gracefully // rather than fail when the telegram channel hasn't been compiled. if !available_channel_names().contains(&"telegram") {